diff --git a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java index 3cd1e2a6e3..38f37cf35d 100644 --- a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java +++ b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java @@ -17,6 +17,7 @@ package org.thingsboard.server.config; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.WebSocketHandler; @@ -40,11 +41,16 @@ public class WebSocketConfiguration implements WebSocketConfigurer { private final WebSocketHandler wsHandler; + @Value("${server.ws.max_text_message_buffer_size:32768}") + private int maxTextMessageBufferSize; + @Value("${server.ws.max_binary_message_buffer_size:32768}") + private int maxBinaryMessageBufferSize; + @Bean public ServletServerContainerFactoryBean createWebSocketContainer() { ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean(); - container.setMaxTextMessageBufferSize(32768); - container.setMaxBinaryMessageBufferSize(32768); + container.setMaxTextMessageBufferSize(maxTextMessageBufferSize); + container.setMaxBinaryMessageBufferSize(maxBinaryMessageBufferSize); return container; } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index a3b9ba2810..c03c6affe2 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -78,6 +78,7 @@ import org.thingsboard.server.service.edge.rpc.processor.telemetry.TelemetryEdge import org.thingsboard.server.service.edge.rpc.processor.user.UserProcessor; import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService; import org.thingsboard.server.service.executors.GrpcCallbackExecutorService; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import java.util.EnumMap; import java.util.List; @@ -104,6 +105,9 @@ public class EdgeContextComponent { } // services + @Autowired + private TelemetrySubscriptionService tsSubService; + @Autowired private AdminSettingsService adminSettingsService; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java index 3003c046b9..445741fd50 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java @@ -41,6 +41,7 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.HasVersion; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TbResource; @@ -193,6 +194,10 @@ public class EdgeMsgConstructorUtils { ) ); + private static void resetVersion(HasVersion entity) { + entity.setVersion(null); + } + public static AlarmUpdateMsg constructAlarmUpdatedMsg(UpdateMsgType msgType, Alarm alarm) { return AlarmUpdateMsg.newBuilder().setMsgType(msgType) .setEntity(JacksonUtil.toString(alarm)) @@ -205,6 +210,7 @@ public class EdgeMsgConstructorUtils { } public static AssetUpdateMsg constructAssetUpdatedMsg(UpdateMsgType msgType, Asset asset) { + resetVersion(asset); return AssetUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(asset)) .setIdMSB(asset.getUuidId().getMostSignificantBits()) .setIdLSB(asset.getUuidId().getLeastSignificantBits()).build(); @@ -218,6 +224,7 @@ public class EdgeMsgConstructorUtils { } public static AssetProfileUpdateMsg constructAssetProfileUpdatedMsg(UpdateMsgType msgType, AssetProfile assetProfile) { + resetVersion(assetProfile); return AssetProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(assetProfile)) .setIdMSB(assetProfile.getId().getId().getMostSignificantBits()) .setIdLSB(assetProfile.getId().getId().getLeastSignificantBits()).build(); @@ -231,6 +238,7 @@ public class EdgeMsgConstructorUtils { } public static CustomerUpdateMsg constructCustomerUpdatedMsg(UpdateMsgType msgType, Customer customer) { + resetVersion(customer); return CustomerUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(customer)) .setIdMSB(customer.getId().getId().getMostSignificantBits()) .setIdLSB(customer.getId().getId().getLeastSignificantBits()).build(); @@ -244,6 +252,7 @@ public class EdgeMsgConstructorUtils { } public static DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard) { + resetVersion(dashboard); return DashboardUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(dashboard)) .setIdMSB(dashboard.getId().getId().getMostSignificantBits()) .setIdLSB(dashboard.getId().getId().getLeastSignificantBits()).build(); @@ -257,6 +266,7 @@ public class EdgeMsgConstructorUtils { } public static DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device) { + resetVersion(device); return DeviceUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(device)) .setIdMSB(device.getId().getId().getMostSignificantBits()) .setIdLSB(device.getId().getId().getLeastSignificantBits()).build(); @@ -270,10 +280,12 @@ public class EdgeMsgConstructorUtils { } public static DeviceCredentialsUpdateMsg constructDeviceCredentialsUpdatedMsg(DeviceCredentials deviceCredentials) { + resetVersion(deviceCredentials); return DeviceCredentialsUpdateMsg.newBuilder().setEntity(JacksonUtil.toString(deviceCredentials)).build(); } public static DeviceProfileUpdateMsg constructDeviceProfileUpdatedMsg(UpdateMsgType msgType, DeviceProfile deviceProfile, EdgeVersion edgeVersion) { + resetVersion(deviceProfile); String entity = getEntityAndFixLwm2mBootstrapShortServerId(deviceProfile, edgeVersion); return DeviceProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(entity) .setIdMSB(deviceProfile.getId().getId().getMostSignificantBits()) @@ -387,6 +399,7 @@ public class EdgeMsgConstructorUtils { } public static EntityViewUpdateMsg constructEntityViewUpdatedMsg(UpdateMsgType msgType, EntityView entityView) { + resetVersion(entityView); return EntityViewUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(entityView)) .setIdMSB(entityView.getId().getId().getMostSignificantBits()) .setIdLSB(entityView.getId().getId().getLeastSignificantBits()).build(); @@ -486,6 +499,7 @@ public class EdgeMsgConstructorUtils { } public static RelationUpdateMsg constructRelationUpdatedMsg(UpdateMsgType msgType, EntityRelation entityRelation) { + resetVersion(entityRelation); return RelationUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(entityRelation)).build(); } @@ -503,6 +517,7 @@ public class EdgeMsgConstructorUtils { } public static RuleChainUpdateMsg constructRuleChainUpdatedMsg(UpdateMsgType msgType, RuleChain ruleChain, boolean isRoot) { + resetVersion(ruleChain); boolean isTemplateRoot = ruleChain.isRoot(); ruleChain.setRoot(isRoot); RuleChainUpdateMsg result = RuleChainUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(ruleChain)) @@ -520,6 +535,7 @@ public class EdgeMsgConstructorUtils { } public static RuleChainMetadataUpdateMsg constructRuleChainMetadataUpdatedMsg(UpdateMsgType msgType, RuleChainMetaData ruleChainMetaData, EdgeVersion edgeVersion) { + resetVersion(ruleChainMetaData); String metaData = sanitizeMetadataForLegacyEdgeVersion(ruleChainMetaData, edgeVersion); return RuleChainMetadataUpdateMsg.newBuilder() @@ -640,6 +656,7 @@ public class EdgeMsgConstructorUtils { } public static TenantUpdateMsg constructTenantUpdateMsg(UpdateMsgType msgType, Tenant tenant) { + resetVersion(tenant); return TenantUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(tenant)).build(); } @@ -648,6 +665,7 @@ public class EdgeMsgConstructorUtils { } public static UserUpdateMsg constructUserUpdatedMsg(UpdateMsgType msgType, User user) { + resetVersion(user); return UserUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(user)) .setIdMSB(user.getId().getId().getMostSignificantBits()) .setIdLSB(user.getId().getId().getLeastSignificantBits()).build(); @@ -665,6 +683,7 @@ public class EdgeMsgConstructorUtils { } public static WidgetsBundleUpdateMsg constructWidgetsBundleUpdateMsg(UpdateMsgType msgType, WidgetsBundle widgetsBundle, List widgets) { + resetVersion(widgetsBundle); return WidgetsBundleUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetsBundle)) .setWidgets(JacksonUtil.toString(widgets)) .setIdMSB(widgetsBundle.getId().getId().getMostSignificantBits()) @@ -680,6 +699,7 @@ public class EdgeMsgConstructorUtils { } public static WidgetTypeUpdateMsg constructWidgetTypeUpdateMsg(UpdateMsgType msgType, WidgetTypeDetails widgetTypeDetails) { + resetVersion(widgetTypeDetails); return WidgetTypeUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetTypeDetails)) .setIdMSB(widgetTypeDetails.getId().getId().getMostSignificantBits()) .setIdLSB(widgetTypeDetails.getId().getId().getLeastSignificantBits()).build(); @@ -694,6 +714,7 @@ public class EdgeMsgConstructorUtils { } public static CalculatedFieldUpdateMsg constructCalculatedFieldUpdatedMsg(UpdateMsgType msgType, CalculatedField calculatedField) { + resetVersion(calculatedField); return CalculatedFieldUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(calculatedField)) .setIdMSB(calculatedField.getId().getId().getMostSignificantBits()) .setIdLSB(calculatedField.getId().getId().getLeastSignificantBits()).build(); @@ -707,6 +728,7 @@ public class EdgeMsgConstructorUtils { } public static AiModelUpdateMsg constructAiModelUpdatedMsg(UpdateMsgType msgType, AiModel aiModel) { + resetVersion(aiModel); return AiModelUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(aiModel)) .setIdMSB(aiModel.getId().getId().getMostSignificantBits()) .setIdLSB(aiModel.getId().getId().getLeastSignificantBits()).build(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java new file mode 100644 index 0000000000..7d0f14c24e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2026 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; + +import com.google.common.util.concurrent.FutureCallback; +import jakarta.annotation.Nullable; +import lombok.AllArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; + +@Slf4j +@AllArgsConstructor +public class AttributeSaveCallback implements FutureCallback { + + private final TenantId tenantId; + private final EdgeId edgeId; + private final String key; + private final Object value; + + @Override + public void onSuccess(@Nullable Void result) { + log.trace("[{}][{}] Successfully updated attribute [{}] with value [{}]", tenantId, edgeId, key, value); + } + + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to update attribute [{}] with value [{}]", tenantId, edgeId, key, value, t); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index 76070bec53..c5822cc494 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -135,9 +135,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Lazy private EdgeContextComponent ctx; - @Autowired - private TelemetrySubscriptionService tsSubService; - @Autowired private TbClusterService clusterService; @@ -553,14 +550,14 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void save(TenantId tenantId, EdgeId edgeId, String key, long value) { log.debug("[{}][{}] Updating long edge telemetry [{}] [{}]", tenantId, edgeId, key, value); if (persistToTelemetry) { - tsSubService.saveTimeseries(TimeseriesSaveRequest.builder() + ctx.getTsSubService().saveTimeseries(TimeseriesSaveRequest.builder() .tenantId(tenantId) .entityId(edgeId) .entry(new LongDataEntry(key, value)) .callback(new AttributeSaveCallback(tenantId, edgeId, key, value)) .build()); } else { - tsSubService.saveAttributes(AttributesSaveRequest.builder() + ctx.getTsSubService().saveAttributes(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(edgeId) .scope(AttributeScope.SERVER_SCOPE) @@ -573,14 +570,14 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void save(TenantId tenantId, EdgeId edgeId, String key, boolean value) { log.debug("[{}][{}] Updating boolean edge telemetry [{}] [{}]", tenantId, edgeId, key, value); if (persistToTelemetry) { - tsSubService.saveTimeseries(TimeseriesSaveRequest.builder() + ctx.getTsSubService().saveTimeseries(TimeseriesSaveRequest.builder() .tenantId(tenantId) .entityId(edgeId) .entry(new BooleanDataEntry(key, value)) .callback(new AttributeSaveCallback(tenantId, edgeId, key, value)) .build()); } else { - tsSubService.saveAttributes(AttributesSaveRequest.builder() + ctx.getTsSubService().saveAttributes(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(edgeId) .scope(AttributeScope.SERVER_SCOPE) @@ -590,32 +587,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private static class AttributeSaveCallback implements FutureCallback { - - private final TenantId tenantId; - private final EdgeId edgeId; - private final String key; - private final Object value; - - AttributeSaveCallback(TenantId tenantId, EdgeId edgeId, String key, Object value) { - this.tenantId = tenantId; - this.edgeId = edgeId; - this.key = key; - this.value = value; - } - - @Override - public void onSuccess(@Nullable Void result) { - log.trace("[{}][{}] Successfully updated attribute [{}] with value [{}]", tenantId, edgeId, key, value); - } - - @Override - public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to update attribute [{}] with value [{}]", tenantId, edgeId, key, value, t); - } - - } - private void pushRuleEngineMessage(TenantId tenantId, Edge edge, long ts, TbMsgType msgType) { try { EdgeId edgeId = edge.getId(); 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 2b026f207f..474a6860d4 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 @@ -25,6 +25,7 @@ import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.data.util.Pair; +import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EdgeUtils; @@ -37,6 +38,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.AttributesSaveResult; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.limit.LimitedApi; @@ -195,7 +197,7 @@ public abstract class EdgeGrpcSession implements Closeable { } startSyncProcess(fullSync); } else { - syncInProgress = false; + updateSyncInProgress(false); } } if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) { @@ -241,6 +243,11 @@ public abstract class EdgeGrpcSession implements Closeable { log.debug("[{}] onConfigurationUpdate [{}]", sessionId, edge); this.tenantId = edge.getTenantId(); this.edge = edge; + if (!this.edge.getCustomerId().equals(edge.getCustomerId())) { + // do not send edge configuration message on customer update + // message send by separate flow from assign_to or unassing_from customer + return; + } EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() .setConfiguration(EdgeMsgConstructorUtils.constructEdgeConfiguration(edge)).build(); ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder() @@ -252,7 +259,7 @@ public abstract class EdgeGrpcSession implements Closeable { public void startSyncProcess(boolean fullSync) { if (!syncInProgress) { log.info("[{}][{}][{}] Staring edge sync process", tenantId, edge.getId(), sessionId); - syncInProgress = true; + updateSyncInProgress(true); interruptGeneralProcessingOnSync(); doSync(new EdgeSyncCursor(ctx, edge, fullSync)); } else { @@ -398,6 +405,18 @@ public abstract class EdgeGrpcSession implements Closeable { ctx.getAttributesService().save(tenantId, edge.getId(), AttributeScope.SERVER_SCOPE, attributeKvEntry); } + private void updateSyncInProgress(Boolean value) { + this.syncInProgress = value; + + ctx.getTsSubService().saveAttributes(AttributesSaveRequest.builder() + .tenantId(tenantId) + .entityId(edge.getId()) + .scope(AttributeScope.SERVER_SCOPE) + .entry(new BooleanDataEntry(DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) + .callback(new AttributeSaveCallback(tenantId, edge.getId(), DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY, value)) + .build()); + } + private void interruptGeneralProcessingOnSync() { log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", tenantId, edge.getId(), sessionId); stopCurrentSendDownlinkMsgsTask(true); @@ -766,7 +785,7 @@ public abstract class EdgeGrpcSession implements Closeable { } private void markSyncCompletedSendEdgeEventUpdate() { - syncInProgress = false; + updateSyncInProgress(false); ctx.getClusterService().onEdgeEventUpdate(new EdgeEventUpdateMsg(edge.getTenantId(), edge.getId())); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java index 55f61ed54f..e9134a0b60 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java @@ -19,6 +19,7 @@ import lombok.Getter; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; 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; @@ -63,7 +64,8 @@ public class EdgeSyncCursor { fetchers.add(new TenantEdgeEventFetcher(ctx.getTenantService())); fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); - fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService())); + fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), TenantId.SYS_TENANT_ID)); + fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService(), edge.getTenantId())); fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService())); fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService())); fetchers.add(new SystemWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService())); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java index 7e4b463573..47e004ac82 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java @@ -17,50 +17,32 @@ package org.thingsboard.server.service.edge.rpc.fetch; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.AdminSettings; 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.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; -import org.thingsboard.server.common.data.id.EdgeId; 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.settings.AdminSettingsService; -import java.util.ArrayList; -import java.util.List; - @AllArgsConstructor @Slf4j -public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { +public class AdminSettingsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final AdminSettingsService adminSettingsService; + private final TenantId fetcherTenantId; @Override - public PageLink getPageLink(int pageSize) { - return null; - } - - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - List result = fetchAdminSettingsForKeys(tenantId, edge.getId(), List.of("general", "mail", "connectivity", "jwt")); - - // return PageData object to be in sync with other fetchers - return new PageData<>(result, 1, result.size(), false); + PageData fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) { + return adminSettingsService.findAllByTenantId(fetcherTenantId, pageLink); } - private List fetchAdminSettingsForKeys(TenantId tenantId, EdgeId edgeId, List keys) { - List result = new ArrayList<>(); - for (String key : keys) { - AdminSettings adminSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, key); - if (adminSettings != null) { - result.add(EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.ADMIN_SETTINGS, - EdgeEventActionType.UPDATED, null, JacksonUtil.valueToTree(adminSettings))); - } - } - return result; + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, AdminSettings adminSettings) { + return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, + EdgeEventActionType.UPDATED, adminSettings.getId(), null); } - } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java index b7337173cd..6c65d141d8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java @@ -41,7 +41,7 @@ public class OAuth2EdgeEventFetcher extends BasePageableEdgeEventFetcher assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class); Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent()); AssetProfileUpdateMsg assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get(); @@ -109,7 +109,7 @@ public class AssetEdgeTest extends AbstractEdgeTest { assetMsg = JacksonUtil.fromString(assetUpdateMsg.getEntity(), Asset.class, true); Assert.assertNotNull(assetMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); - Assert.assertEquals(savedAsset, assetMsg); + compareHasVersionEntities(savedAsset, assetMsg); assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class); Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent()); assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get(); diff --git a/application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java index 68548558f7..210b2799e0 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java @@ -53,7 +53,7 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest { AssetProfileUpdateMsg assetProfileUpdateMsg = (AssetProfileUpdateMsg) latestMessage; AssetProfile assetProfileMsg = JacksonUtil.fromString(assetProfileUpdateMsg.getEntity(), AssetProfile.class, true); Assert.assertNotNull(assetProfileMsg); - Assert.assertEquals(assetProfile, assetProfileMsg); + compareHasVersionEntities(assetProfile, assetProfileMsg); Assert.assertEquals("Building", assetProfileMsg.getName()); Assert.assertEquals(buildingsRuleChainId, assetProfileMsg.getDefaultEdgeRuleChainId()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetProfileUpdateMsg.getMsgType()); diff --git a/application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java index 6199657936..60701cee1a 100644 --- a/application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java @@ -147,7 +147,7 @@ public class CalculatedFieldEdgeTest extends AbstractEdgeTest { CalculatedFieldUpdateMsg calculatedFieldUpdateMsg = (CalculatedFieldUpdateMsg) latestMessage; CalculatedField calculatedFieldFromEdge = JacksonUtil.fromString(calculatedFieldUpdateMsg.getEntity(), CalculatedField.class, true); Assert.assertNotNull(calculatedFieldFromEdge); - Assert.assertEquals(savedCalculatedField, calculatedFieldFromEdge); + compareHasVersionEntities(savedCalculatedField, calculatedFieldFromEdge); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, calculatedFieldUpdateMsg.getMsgType()); } diff --git a/application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java index f3d12675a3..578175dc46 100644 --- a/application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java @@ -60,7 +60,7 @@ public class CustomerEdgeTest extends AbstractEdgeTest { CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); - Assert.assertEquals(savedCustomer, customerMsg); + compareHasVersionEntities(savedCustomer, customerMsg); testAutoGeneratedCodeByProtobuf(customerUpdateMsg); // update customer @@ -73,7 +73,7 @@ public class CustomerEdgeTest extends AbstractEdgeTest { customerUpdateMsg = (CustomerUpdateMsg) latestMessage; customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); - Assert.assertEquals(savedCustomer, customerMsg); + compareHasVersionEntities(savedCustomer, customerMsg); // delete customer edgeImitator.expectMessageAmount(2); 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 9c70db267e..788dff02c0 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java @@ -201,7 +201,7 @@ public class DashboardEdgeTest extends AbstractEdgeTest { CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); - Assert.assertEquals(savedCustomer, customerMsg); + compareHasVersionEntities(savedCustomer, customerMsg); Dashboard dashboard = buildDashboardForUplinkMsg(savedCustomer); 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 9349f79c02..369cda37aa 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java @@ -141,7 +141,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest { Device deviceFromMsg = JacksonUtil.fromString(deviceUpdateMsg.getEntity(), Device.class, true); Assert.assertNotNull(deviceFromMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); - Assert.assertEquals(savedDevice, deviceFromMsg); + compareHasVersionEntities(savedDevice, deviceFromMsg); Assert.assertEquals(savedDevice.getId(), deviceFromMsg.getId()); Assert.assertEquals(savedDevice.getName(), deviceFromMsg.getName()); Assert.assertEquals(savedDevice.getType(), deviceFromMsg.getType()); @@ -222,7 +222,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest { Assert.assertTrue(latestMessage instanceof DeviceCredentialsUpdateMsg); DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = (DeviceCredentialsUpdateMsg) latestMessage; DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true); - Assert.assertEquals(deviceCredentials, deviceCredentialsMsg); + compareHasVersionEntities(deviceCredentials, deviceCredentialsMsg); // update device credentials - X509_CERTIFICATE edgeImitator.expectMessageAmount(1); @@ -272,7 +272,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest { Device deviceMsg = JacksonUtil.fromString(deviceUpdateMsg.getEntity(), Device.class, true); Assert.assertNotNull(deviceMsg); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); - Assert.assertEquals(savedDevice, deviceMsg); + compareHasVersionEntities(savedDevice, deviceMsg); Assert.assertEquals(firmwareOtaPackageInfo.getId(), deviceMsg.getFirmwareId()); Assert.assertEquals(softwareOtaPackageInfo.getId(), deviceMsg.getSoftwareId()); deviceData = deviceMsg.getDeviceData(); @@ -387,7 +387,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest { DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true); Assert.assertNotNull(deviceCredentialsMsg); Assert.assertEquals(device.getId(), deviceCredentialsMsg.getDeviceId()); - Assert.assertEquals(deviceCredentials, deviceCredentialsMsg); + compareHasVersionEntities(deviceCredentials, deviceCredentialsMsg); } @Test diff --git a/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java index f44b4e295e..3ac7b3d648 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java @@ -83,7 +83,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); // update device profile @@ -108,7 +108,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); // delete profile edgeImitator.expectMessageAmount(1); @@ -146,7 +146,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); // delete profile when edge is offline @@ -186,7 +186,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(DeviceTransportType.SNMP, deviceProfileMsg.getTransportType()); @@ -224,7 +224,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(DeviceTransportType.LWM2M, deviceProfileMsg.getTransportType()); @@ -273,7 +273,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); Assert.assertNotNull(deviceProfileMsg); - Assert.assertEquals(deviceProfile, deviceProfileMsg); + compareHasVersionEntities(deviceProfile, deviceProfileMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(DeviceTransportType.COAP, deviceProfileMsg.getTransportType()); diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeTest.java index 334680f7a7..e9dbf6d443 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/EdgeTest.java @@ -15,10 +15,12 @@ */ package org.thingsboard.server.edge; +import com.fasterxml.jackson.databind.JsonNode; import org.junit.Assert; import org.junit.Test; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; @@ -26,10 +28,14 @@ import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.gen.edge.v1.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; +import java.util.List; import java.util.Optional; import java.util.UUID; +import static org.assertj.core.api.Assertions.assertThat; + @DaoSqlTest public class EdgeTest extends AbstractEdgeTest { @@ -57,7 +63,7 @@ public class EdgeTest extends AbstractEdgeTest { CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); - Assert.assertEquals(savedCustomer, customerMsg); + compareHasVersionEntities(savedCustomer, customerMsg); // unassign edge from customer edgeImitator.expectMessageAmount(2); @@ -76,4 +82,23 @@ public class EdgeTest extends AbstractEdgeTest { Assert.assertEquals(savedCustomer.getUuidId().getMostSignificantBits(), customerUpdateMsg.getIdMSB()); Assert.assertEquals(savedCustomer.getUuidId().getLeastSignificantBits(), customerUpdateMsg.getIdLSB()); } + + @Test + public void testSyncEdge_attributeUpdated() throws Exception { + getWsClient().subscribeForAttributes(edge.getId(), TbAttributeSubscriptionScope.SERVER_SCOPE.name(), List.of(DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY)); + + doPost("/api/edge/sync/" + edge.getId()); + + // wait for sync to start + waitForEdgeSyncInProgressEqualsValue(true); + + // wait for sync to end + waitForEdgeSyncInProgressEqualsValue(false); + } + + private void waitForEdgeSyncInProgressEqualsValue(Boolean value) { + getWsClient().registerWaitForUpdate(); + JsonNode update = JacksonUtil.toJsonNode(getWsClient().waitForUpdate()); + assertThat(update.get("data").get(DataConstants.EDGE_SYNC_IN_PROGRESS_ATTR_KEY).get(0).get(1).asBoolean()).isEqualTo(value); + } } diff --git a/application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java index ce7b915457..640545c5d4 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java @@ -68,7 +68,7 @@ public class EntityViewEdgeTest extends AbstractEdgeTest { EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage; EntityView entityView = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true); Assert.assertNotNull(entityView); - Assert.assertEquals(savedEntityView, entityView); + compareHasVersionEntities(savedEntityView, entityView); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType()); // request entity view(s) for device @@ -265,7 +265,7 @@ public class EntityViewEdgeTest extends AbstractEdgeTest { EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage; EntityView entityViewMsg = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true); Assert.assertNotNull(entityViewMsg); - Assert.assertEquals(entityView, entityViewMsg); + compareHasVersionEntities(entityView, entityViewMsg); Assert.assertEquals(device.getId(), entityViewMsg.getEntityId()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType()); testAutoGeneratedCodeByProtobuf(entityViewUpdateMsg); diff --git a/application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java index 054f67585e..cc23440deb 100644 --- a/application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java @@ -57,7 +57,7 @@ public class RelationEdgeTest extends AbstractEdgeTest { RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); Assert.assertNotNull(entityRelation); - Assert.assertEquals(relation, entityRelation); + compareHasVersionEntities(relation, entityRelation); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); // delete relation @@ -76,7 +76,7 @@ public class RelationEdgeTest extends AbstractEdgeTest { relationUpdateMsg = (RelationUpdateMsg) latestMessage; entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); Assert.assertNotNull(entityRelation); - Assert.assertEquals(deletedRelation, entityRelation); + compareHasVersionEntities(deletedRelation, entityRelation); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); } @@ -155,7 +155,7 @@ public class RelationEdgeTest extends AbstractEdgeTest { RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); Assert.assertNotNull(entityRelation); - Assert.assertEquals(deviceToAssetRelation, entityRelation); + compareHasVersionEntities(deviceToAssetRelation, entityRelation); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); } @@ -177,7 +177,7 @@ public class RelationEdgeTest extends AbstractEdgeTest { RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); Assert.assertNotNull(entityRelation); - Assert.assertEquals(relation, entityRelation); + compareHasVersionEntities(relation, entityRelation); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); // delete relation @@ -196,7 +196,7 @@ public class RelationEdgeTest extends AbstractEdgeTest { relationUpdateMsg = (RelationUpdateMsg) latestMessage; entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); Assert.assertNotNull(entityRelation); - Assert.assertEquals(deletedRelation, entityRelation); + compareHasVersionEntities(deletedRelation, entityRelation); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); } diff --git a/application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java index 283e8bb66b..812357ac7d 100644 --- a/application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java @@ -50,7 +50,7 @@ public class TenantEdgeTest extends AbstractEdgeTest { TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileUpdateMsgOpt.get(); Tenant tenantMsg = JacksonUtil.fromString(tenantUpdateMsg.getEntity(), Tenant.class, true); Assert.assertNotNull(tenantMsg); - Assert.assertEquals(savedTenant, tenantMsg); + compareHasVersionEntities(savedTenant, tenantMsg); TenantProfile tenantProfileMsg = JacksonUtil.fromString(tenantProfileUpdateMsg.getEntity(), TenantProfile.class, true); Assert.assertNotNull(tenantProfileMsg); Assert.assertEquals(tenantMsg.getTenantProfileId(), tenantProfileMsg.getId()); @@ -75,7 +75,7 @@ public class TenantEdgeTest extends AbstractEdgeTest { Assert.assertNotNull(tenantProfileMsg); // tenant update Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, tenantUpdateMsg.getMsgType()); - Assert.assertEquals(savedTenant, tenantMsg); + compareHasVersionEntities(savedTenant, tenantMsg); Assert.assertEquals(savedTenant.getTenantProfileId(), tenantProfileMsg.getId()); } diff --git a/application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java index d6afa81262..ddaa867aee 100644 --- a/application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java @@ -68,7 +68,7 @@ public class WidgetEdgeTest extends AbstractEdgeTest { WidgetType widgetsType = JacksonUtil.fromString(widgetTypeUpdateMsg.getEntity(), WidgetType.class, true); Assert.assertNotNull(widgetsType); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, widgetTypeUpdateMsg.getMsgType()); - Assert.assertEquals(savedWidgetType, widgetsType); + compareHasVersionEntities(savedWidgetType, widgetsType); // update widget bundle edgeImitator.expectMessageAmount(1); diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java index 4bbb51068c..3bf5fa2f25 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java @@ -38,12 +38,30 @@ import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNode; import org.thingsboard.rule.engine.telemetry.TbCalculatedFieldsNode; import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; +import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.Dashboard; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.asset.AssetProfile; 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.AssetId; +import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DashboardId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; +import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.gen.edge.v1.EdgeVersion; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; @@ -278,4 +296,167 @@ public class EdgeMsgConstructorUtilsTest { edgeEvent.setBody(body); return edgeEvent; } + + @Test + public void testConstructAssetUpdatedMsg_versionIsReset() { + Asset asset = new Asset(); + asset.setId(new AssetId(UUID.randomUUID())); + asset.setName("Test Asset"); + asset.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructAssetUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, asset).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "Asset version should be null in serialized message"); + } + + @Test + public void testConstructAssetProfileUpdatedMsg_versionIsReset() { + AssetProfile assetProfile = new AssetProfile(); + assetProfile.setId(new AssetProfileId(UUID.randomUUID())); + assetProfile.setName("Test Asset Profile"); + assetProfile.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructAssetProfileUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, assetProfile).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "AssetProfile version should be null in serialized message"); + } + + @Test + public void testConstructCustomerUpdatedMsg_versionIsReset() { + Customer customer = new Customer(); + customer.setId(new CustomerId(UUID.randomUUID())); + customer.setTitle("Test Customer"); + customer.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructCustomerUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, customer).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "Customer version should be null in serialized message"); + } + + @Test + public void testConstructDashboardUpdatedMsg_versionIsReset() { + Dashboard dashboard = new Dashboard(); + dashboard.setId(new DashboardId(UUID.randomUUID())); + dashboard.setTitle("Test Dashboard"); + dashboard.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructDashboardUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, dashboard).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "Dashboard version should be null in serialized message"); + } + + @Test + public void testConstructDeviceUpdatedMsg_versionIsReset() { + Device device = new Device(); + device.setId(new DeviceId(UUID.randomUUID())); + device.setName("Test Device"); + device.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructDeviceUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, device).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "Device version should be null in serialized message"); + } + + @Test + public void testConstructDeviceCredentialsUpdatedMsg_versionIsReset() { + DeviceCredentials credentials = new DeviceCredentials(); + credentials.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructDeviceCredentialsUpdatedMsg(credentials).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "DeviceCredentials version should be null in serialized message"); + } + + @Test + public void testConstructEntityViewUpdatedMsg_versionIsReset() { + EntityView entityView = new EntityView(); + entityView.setId(new EntityViewId(UUID.randomUUID())); + entityView.setName("Test EntityView"); + entityView.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructEntityViewUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, entityView).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "EntityView version should be null in serialized message"); + } + + @Test + public void testConstructRelationUpdatedMsg_versionIsReset() { + EntityRelation relation = new EntityRelation(); + relation.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructRelationUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, relation).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "EntityRelation version should be null in serialized message"); + } + + @Test + public void testConstructRuleChainUpdatedMsg_versionIsReset() { + RuleChain ruleChain = new RuleChain(); + ruleChain.setId(new org.thingsboard.server.common.data.id.RuleChainId(UUID.randomUUID())); + ruleChain.setName("Test RuleChain"); + ruleChain.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructRuleChainUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, ruleChain, false).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "RuleChain version should be null in serialized message"); + } + + @Test + public void testConstructTenantUpdateMsg_versionIsReset() { + Tenant tenant = new Tenant(); + tenant.setId(TenantId.fromUUID(UUID.randomUUID())); + tenant.setTitle("Test Tenant"); + tenant.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructTenantUpdateMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, tenant).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "Tenant version should be null in serialized message"); + } + + @Test + public void testConstructUserUpdatedMsg_versionIsReset() { + User user = new User(); + user.setId(new UserId(UUID.randomUUID())); + user.setEmail("test@test.com"); + user.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructUserUpdatedMsg(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, user).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "User version should be null in serialized message"); + } + + @Test + public void testConstructRuleChainMetadataUpdatedMsg_versionIsReset() { + RuleChainMetaData metaData = new RuleChainMetaData(); + metaData.setVersion(42L); + + String entity = EdgeMsgConstructorUtils.constructRuleChainMetadataUpdatedMsg( + UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, metaData, EdgeVersion.V_4_0_0).getEntity(); + JsonNode json = JacksonUtil.toJsonNode(entity); + + Assertions.assertTrue(json.get("version") == null || json.get("version").isNull(), + "RuleChainMetaData version should be null in serialized message"); + } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsService.java index dc7af6275d..ce243853e8 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsService.java @@ -18,6 +18,8 @@ package org.thingsboard.server.dao.settings; import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.id.AdminSettingsId; 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.entity.EntityDaoService; public interface AdminSettingsService extends EntityDaoService { @@ -28,6 +30,8 @@ public interface AdminSettingsService extends EntityDaoService { AdminSettings findAdminSettingsByTenantIdAndKey(TenantId tenantId, String key); + PageData findAllByTenantId(TenantId tenantId, PageLink pageLink); + AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings); boolean deleteAdminSettingsByTenantIdAndKey(TenantId tenantId, String key); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index de46575e06..04f3143834 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -127,6 +127,7 @@ public class DataConstants { public static final String EDGE_MSG_SOURCE = "edge"; public static final String MSG_SOURCE_KEY = "source"; public static final String EDGE_VERSION_ATTR_KEY = "edgeVersion"; + public static final String EDGE_SYNC_IN_PROGRESS_ATTR_KEY = "syncInProgress"; public static final String LAST_CONNECTED_GATEWAY = "lastConnectedGateway"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java index b07c159fa9..834e19a1c2 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java @@ -38,7 +38,7 @@ public enum EdgeEventType { TENANT_PROFILE(true, EntityType.TENANT_PROFILE), WIDGETS_BUNDLE(true, EntityType.WIDGETS_BUNDLE), WIDGET_TYPE(true, EntityType.WIDGET_TYPE), - ADMIN_SETTINGS(true, null), + ADMIN_SETTINGS(true, EntityType.ADMIN_SETTINGS), OTA_PACKAGE(true, EntityType.OTA_PACKAGE), QUEUE(true, EntityType.QUEUE), NOTIFICATION_RULE(true, EntityType.NOTIFICATION_RULE), diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java index 93851fa813..7baa8e72f2 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java @@ -114,6 +114,7 @@ public class EntityIdFactory { case DOMAIN -> new DomainId(uuid); case CALCULATED_FIELD -> new CalculatedFieldId(uuid); case AI_MODEL -> new AiModelId(uuid); + case ADMIN_SETTINGS -> new AdminSettingsId(uuid); default -> throw new IllegalArgumentException("EdgeEventType " + edgeEventType + " is not supported!"); }; } diff --git a/common/util/src/test/java/org/thingsboard/common/util/SsrfProtectionValidatorTest.java b/common/util/src/test/java/org/thingsboard/common/util/SsrfProtectionValidatorTest.java index ec6e51db6d..6cb2d21a9a 100644 --- a/common/util/src/test/java/org/thingsboard/common/util/SsrfProtectionValidatorTest.java +++ b/common/util/src/test/java/org/thingsboard/common/util/SsrfProtectionValidatorTest.java @@ -27,12 +27,9 @@ import java.util.List; import static org.assertj.core.api.Assertions.assertThatNoException; import static org.assertj.core.api.Assertions.assertThatThrownBy; +@ResourceLock("SsrfProtectionValidatorTest") // some tests mutate static additional-blocked-hosts public class SsrfProtectionValidatorTest { - // JUnit 5 @ResourceLock ensures that tests modifying SsrfProtectionValidator's static - // additional blocked hosts never run concurrently with each other (parallel execution is enabled). - private static final String SYNC_LOCK = "SsrfProtectionValidatorTest"; - @ParameterizedTest @ValueSource(strings = { "http://example.com", @@ -207,7 +204,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedSingleIp() { try { SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("8.8.8.8")); @@ -222,7 +218,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedCidrSlash10() { try { // Use 44.0.0.0/10 (not blocked by default) to verify CIDR /10 matching @@ -243,7 +238,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedCidrSlash24() { try { SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("198.51.100.0/24")); @@ -261,7 +255,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedHostnameViaValidateUri() { try { SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("evil.corp")); @@ -275,7 +268,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedHostnameCaseInsensitive() { try { SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("My-Service.Corp")); @@ -289,7 +281,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testSetAdditionalBlockedHostsEmptyAndNull() { // Should not throw SsrfProtectionValidator.setAdditionalBlockedHosts(Collections.emptyList()); @@ -307,7 +298,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedCidrViaValidateUri() { // 203.0.113.0/24 (TEST-NET-3) is not blocked by default URI uri = URI.create("http://203.0.113.1"); @@ -323,7 +313,6 @@ public class SsrfProtectionValidatorTest { } @Test - @ResourceLock(SYNC_LOCK) void testAdditionalBlockedMixedConfig() { try { SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("203.0.113.0/24", "evil.corp", "8.8.8.8")); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index fa955fa72e..70f5e94d92 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -66,7 +66,6 @@ import org.thingsboard.server.dao.entity.EntityCountService; import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; -import org.thingsboard.server.exception.DataValidationException; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; @@ -75,6 +74,7 @@ import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.sql.JpaExecutorService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.user.UserService; +import org.thingsboard.server.exception.DataValidationException; import java.util.ArrayList; import java.util.Collections; @@ -239,9 +239,10 @@ public class EdgeServiceImpl extends AbstractCachedEntityService { AdminSettings findByTenantIdAndKey(UUID tenantId, String key); + PageData findAllByTenantId(TenantId tenantId, PageLink pageLink); + boolean removeByTenantIdAndKey(UUID tenantId, String key); void removeByTenantId(UUID tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java index d6708aabfe..71f0eb268e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java @@ -20,6 +20,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.EntityType; @@ -27,6 +28,9 @@ import org.thingsboard.server.common.data.id.AdminSettingsId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; 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.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.Validator; @@ -44,6 +48,9 @@ public class AdminSettingsServiceImpl implements AdminSettingsService { @Autowired private DataValidator adminSettingsValidator; + @Autowired + protected ApplicationEventPublisher eventPublisher; + @Override public AdminSettings findAdminSettingsById(TenantId tenantId, AdminSettingsId adminSettingsId) { log.trace("Executing findAdminSettingsById [{}]", adminSettingsId); @@ -63,10 +70,15 @@ public class AdminSettingsServiceImpl implements AdminSettingsService { return adminSettingsDao.findByTenantIdAndKey(tenantId.getId(), key); } + @Override + public PageData findAllByTenantId(TenantId tenantId, PageLink pageLink) { + return adminSettingsDao.findAllByTenantId(tenantId, pageLink); + } + @Override public AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings) { log.trace("Executing saveAdminSettings [{}]", adminSettings); - adminSettingsValidator.validate(adminSettings, data -> tenantId); + AdminSettings oldAdminSettings = adminSettingsValidator.validate(adminSettings, data -> tenantId); if (adminSettings.getKey().equals("mail")) { AdminSettings mailSettings = findAdminSettingsByKey(tenantId, "mail"); if (mailSettings != null) { @@ -84,7 +96,10 @@ public class AdminSettingsServiceImpl implements AdminSettingsService { if (adminSettings.getTenantId() == null) { adminSettings.setTenantId(TenantId.SYS_TENANT_ID); } - return adminSettingsDao.save(tenantId, adminSettings); + AdminSettings savedAdminSettings = adminSettingsDao.save(tenantId, adminSettings); + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(savedAdminSettings.getTenantId()).entityId(savedAdminSettings.getId()) + .entity(savedAdminSettings).oldEntity(oldAdminSettings).created(adminSettings.getId() == null).build()); + return savedAdminSettings; } @Override diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AdminSettingsServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AdminSettingsServiceTest.java index 8423599c90..50d0b852f1 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/AdminSettingsServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/AdminSettingsServiceTest.java @@ -25,9 +25,12 @@ import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.exception.DataValidationException; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.settings.AdminSettingsService; +import org.thingsboard.server.exception.DataValidationException; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; @@ -113,4 +116,24 @@ public class AdminSettingsServiceTest extends AbstractServiceTest { }).hasMessageContaining("already exists"); } + @Test + public void testFindAllByTenantId() { + int pageSize = 10; + int totalElements = 100; + + for (int i = 0; i < totalElements; i++) { + AdminSettings settings = new AdminSettings(); + settings.setTenantId(tenantId); + String key = RandomStringUtils.randomAlphanumeric(15); + settings.setKey(key); + settings.setJsonValue(JacksonUtil.newObjectNode().put("value", i)); + adminSettingsService.saveAdminSettings(tenantId, settings); + } + + PageData pageData = adminSettingsService.findAllByTenantId(tenantId, new PageLink(pageSize)); + assertThat(pageData.getData().size()).isEqualTo(pageSize); + assertThat(pageData.getTotalElements()).isEqualTo(totalElements); + + } + }