Browse Source

Merge pull request #15158 from thingsboard/rc

Merge rc into master
pull/15236/head
Viacheslav Klimov 7 months ago
committed by GitHub
parent
commit
16084f77f2
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 10
      application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  3. 22
      application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java
  4. 43
      application/src/main/java/org/thingsboard/server/service/edge/rpc/AttributeSaveCallback.java
  5. 37
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  6. 25
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  7. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  8. 34
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java
  10. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  11. 10
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/settings/AdminSettingsEdgeProcessor.java
  12. 4
      application/src/main/resources/thingsboard.yml
  13. 9
      application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
  14. 79
      application/src/test/java/org/thingsboard/server/edge/AdminSettingsEdgeTest.java
  15. 4
      application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java
  16. 2
      application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java
  17. 2
      application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java
  18. 4
      application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java
  19. 2
      application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java
  20. 8
      application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java
  21. 12
      application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java
  22. 27
      application/src/test/java/org/thingsboard/server/edge/EdgeTest.java
  23. 4
      application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java
  24. 10
      application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java
  25. 4
      application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java
  26. 2
      application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java
  27. 181
      application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java
  28. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsService.java
  29. 1
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  30. 2
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java
  31. 1
      common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java
  32. 13
      common/util/src/test/java/org/thingsboard/common/util/SsrfProtectionValidatorTest.java
  33. 5
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  34. 4
      dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsDao.java
  35. 19
      dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java
  36. 25
      dao/src/test/java/org/thingsboard/server/dao/service/AdminSettingsServiceTest.java

10
application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java

@ -17,6 +17,7 @@ package org.thingsboard.server.config;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.WebSocketHandler; import org.springframework.web.socket.WebSocketHandler;
@ -40,11 +41,16 @@ public class WebSocketConfiguration implements WebSocketConfigurer {
private final WebSocketHandler wsHandler; 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 @Bean
public ServletServerContainerFactoryBean createWebSocketContainer() { public ServletServerContainerFactoryBean createWebSocketContainer() {
ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean(); ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean();
container.setMaxTextMessageBufferSize(32768); container.setMaxTextMessageBufferSize(maxTextMessageBufferSize);
container.setMaxBinaryMessageBufferSize(32768); container.setMaxBinaryMessageBufferSize(maxBinaryMessageBufferSize);
return container; return container;
} }

4
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.processor.user.UserProcessor;
import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService; import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService;
import org.thingsboard.server.service.executors.GrpcCallbackExecutorService; import org.thingsboard.server.service.executors.GrpcCallbackExecutorService;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import java.util.EnumMap; import java.util.EnumMap;
import java.util.List; import java.util.List;
@ -104,6 +105,9 @@ public class EdgeContextComponent {
} }
// services // services
@Autowired
private TelemetrySubscriptionService tsSubService;
@Autowired @Autowired
private AdminSettingsService adminSettingsService; private AdminSettingsService adminSettingsService;

22
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.Device;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityView; 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.OtaPackage;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TbResource; 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) { public static AlarmUpdateMsg constructAlarmUpdatedMsg(UpdateMsgType msgType, Alarm alarm) {
return AlarmUpdateMsg.newBuilder().setMsgType(msgType) return AlarmUpdateMsg.newBuilder().setMsgType(msgType)
.setEntity(JacksonUtil.toString(alarm)) .setEntity(JacksonUtil.toString(alarm))
@ -205,6 +210,7 @@ public class EdgeMsgConstructorUtils {
} }
public static AssetUpdateMsg constructAssetUpdatedMsg(UpdateMsgType msgType, Asset asset) { public static AssetUpdateMsg constructAssetUpdatedMsg(UpdateMsgType msgType, Asset asset) {
resetVersion(asset);
return AssetUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(asset)) return AssetUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(asset))
.setIdMSB(asset.getUuidId().getMostSignificantBits()) .setIdMSB(asset.getUuidId().getMostSignificantBits())
.setIdLSB(asset.getUuidId().getLeastSignificantBits()).build(); .setIdLSB(asset.getUuidId().getLeastSignificantBits()).build();
@ -218,6 +224,7 @@ public class EdgeMsgConstructorUtils {
} }
public static AssetProfileUpdateMsg constructAssetProfileUpdatedMsg(UpdateMsgType msgType, AssetProfile assetProfile) { public static AssetProfileUpdateMsg constructAssetProfileUpdatedMsg(UpdateMsgType msgType, AssetProfile assetProfile) {
resetVersion(assetProfile);
return AssetProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(assetProfile)) return AssetProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(assetProfile))
.setIdMSB(assetProfile.getId().getId().getMostSignificantBits()) .setIdMSB(assetProfile.getId().getId().getMostSignificantBits())
.setIdLSB(assetProfile.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(assetProfile.getId().getId().getLeastSignificantBits()).build();
@ -231,6 +238,7 @@ public class EdgeMsgConstructorUtils {
} }
public static CustomerUpdateMsg constructCustomerUpdatedMsg(UpdateMsgType msgType, Customer customer) { public static CustomerUpdateMsg constructCustomerUpdatedMsg(UpdateMsgType msgType, Customer customer) {
resetVersion(customer);
return CustomerUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(customer)) return CustomerUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(customer))
.setIdMSB(customer.getId().getId().getMostSignificantBits()) .setIdMSB(customer.getId().getId().getMostSignificantBits())
.setIdLSB(customer.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(customer.getId().getId().getLeastSignificantBits()).build();
@ -244,6 +252,7 @@ public class EdgeMsgConstructorUtils {
} }
public static DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard) { public static DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard) {
resetVersion(dashboard);
return DashboardUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(dashboard)) return DashboardUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(dashboard))
.setIdMSB(dashboard.getId().getId().getMostSignificantBits()) .setIdMSB(dashboard.getId().getId().getMostSignificantBits())
.setIdLSB(dashboard.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(dashboard.getId().getId().getLeastSignificantBits()).build();
@ -257,6 +266,7 @@ public class EdgeMsgConstructorUtils {
} }
public static DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device) { public static DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device) {
resetVersion(device);
return DeviceUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(device)) return DeviceUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(device))
.setIdMSB(device.getId().getId().getMostSignificantBits()) .setIdMSB(device.getId().getId().getMostSignificantBits())
.setIdLSB(device.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(device.getId().getId().getLeastSignificantBits()).build();
@ -270,10 +280,12 @@ public class EdgeMsgConstructorUtils {
} }
public static DeviceCredentialsUpdateMsg constructDeviceCredentialsUpdatedMsg(DeviceCredentials deviceCredentials) { public static DeviceCredentialsUpdateMsg constructDeviceCredentialsUpdatedMsg(DeviceCredentials deviceCredentials) {
resetVersion(deviceCredentials);
return DeviceCredentialsUpdateMsg.newBuilder().setEntity(JacksonUtil.toString(deviceCredentials)).build(); return DeviceCredentialsUpdateMsg.newBuilder().setEntity(JacksonUtil.toString(deviceCredentials)).build();
} }
public static DeviceProfileUpdateMsg constructDeviceProfileUpdatedMsg(UpdateMsgType msgType, DeviceProfile deviceProfile, EdgeVersion edgeVersion) { public static DeviceProfileUpdateMsg constructDeviceProfileUpdatedMsg(UpdateMsgType msgType, DeviceProfile deviceProfile, EdgeVersion edgeVersion) {
resetVersion(deviceProfile);
String entity = getEntityAndFixLwm2mBootstrapShortServerId(deviceProfile, edgeVersion); String entity = getEntityAndFixLwm2mBootstrapShortServerId(deviceProfile, edgeVersion);
return DeviceProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(entity) return DeviceProfileUpdateMsg.newBuilder().setMsgType(msgType).setEntity(entity)
.setIdMSB(deviceProfile.getId().getId().getMostSignificantBits()) .setIdMSB(deviceProfile.getId().getId().getMostSignificantBits())
@ -387,6 +399,7 @@ public class EdgeMsgConstructorUtils {
} }
public static EntityViewUpdateMsg constructEntityViewUpdatedMsg(UpdateMsgType msgType, EntityView entityView) { public static EntityViewUpdateMsg constructEntityViewUpdatedMsg(UpdateMsgType msgType, EntityView entityView) {
resetVersion(entityView);
return EntityViewUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(entityView)) return EntityViewUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(entityView))
.setIdMSB(entityView.getId().getId().getMostSignificantBits()) .setIdMSB(entityView.getId().getId().getMostSignificantBits())
.setIdLSB(entityView.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(entityView.getId().getId().getLeastSignificantBits()).build();
@ -486,6 +499,7 @@ public class EdgeMsgConstructorUtils {
} }
public static RelationUpdateMsg constructRelationUpdatedMsg(UpdateMsgType msgType, EntityRelation entityRelation) { public static RelationUpdateMsg constructRelationUpdatedMsg(UpdateMsgType msgType, EntityRelation entityRelation) {
resetVersion(entityRelation);
return RelationUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(entityRelation)).build(); 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) { public static RuleChainUpdateMsg constructRuleChainUpdatedMsg(UpdateMsgType msgType, RuleChain ruleChain, boolean isRoot) {
resetVersion(ruleChain);
boolean isTemplateRoot = ruleChain.isRoot(); boolean isTemplateRoot = ruleChain.isRoot();
ruleChain.setRoot(isRoot); ruleChain.setRoot(isRoot);
RuleChainUpdateMsg result = RuleChainUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(ruleChain)) 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) { public static RuleChainMetadataUpdateMsg constructRuleChainMetadataUpdatedMsg(UpdateMsgType msgType, RuleChainMetaData ruleChainMetaData, EdgeVersion edgeVersion) {
resetVersion(ruleChainMetaData);
String metaData = sanitizeMetadataForLegacyEdgeVersion(ruleChainMetaData, edgeVersion); String metaData = sanitizeMetadataForLegacyEdgeVersion(ruleChainMetaData, edgeVersion);
return RuleChainMetadataUpdateMsg.newBuilder() return RuleChainMetadataUpdateMsg.newBuilder()
@ -640,6 +656,7 @@ public class EdgeMsgConstructorUtils {
} }
public static TenantUpdateMsg constructTenantUpdateMsg(UpdateMsgType msgType, Tenant tenant) { public static TenantUpdateMsg constructTenantUpdateMsg(UpdateMsgType msgType, Tenant tenant) {
resetVersion(tenant);
return TenantUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(tenant)).build(); 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) { public static UserUpdateMsg constructUserUpdatedMsg(UpdateMsgType msgType, User user) {
resetVersion(user);
return UserUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(user)) return UserUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(user))
.setIdMSB(user.getId().getId().getMostSignificantBits()) .setIdMSB(user.getId().getId().getMostSignificantBits())
.setIdLSB(user.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(user.getId().getId().getLeastSignificantBits()).build();
@ -665,6 +683,7 @@ public class EdgeMsgConstructorUtils {
} }
public static WidgetsBundleUpdateMsg constructWidgetsBundleUpdateMsg(UpdateMsgType msgType, WidgetsBundle widgetsBundle, List<String> widgets) { public static WidgetsBundleUpdateMsg constructWidgetsBundleUpdateMsg(UpdateMsgType msgType, WidgetsBundle widgetsBundle, List<String> widgets) {
resetVersion(widgetsBundle);
return WidgetsBundleUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetsBundle)) return WidgetsBundleUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetsBundle))
.setWidgets(JacksonUtil.toString(widgets)) .setWidgets(JacksonUtil.toString(widgets))
.setIdMSB(widgetsBundle.getId().getId().getMostSignificantBits()) .setIdMSB(widgetsBundle.getId().getId().getMostSignificantBits())
@ -680,6 +699,7 @@ public class EdgeMsgConstructorUtils {
} }
public static WidgetTypeUpdateMsg constructWidgetTypeUpdateMsg(UpdateMsgType msgType, WidgetTypeDetails widgetTypeDetails) { public static WidgetTypeUpdateMsg constructWidgetTypeUpdateMsg(UpdateMsgType msgType, WidgetTypeDetails widgetTypeDetails) {
resetVersion(widgetTypeDetails);
return WidgetTypeUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetTypeDetails)) return WidgetTypeUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(widgetTypeDetails))
.setIdMSB(widgetTypeDetails.getId().getId().getMostSignificantBits()) .setIdMSB(widgetTypeDetails.getId().getId().getMostSignificantBits())
.setIdLSB(widgetTypeDetails.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(widgetTypeDetails.getId().getId().getLeastSignificantBits()).build();
@ -694,6 +714,7 @@ public class EdgeMsgConstructorUtils {
} }
public static CalculatedFieldUpdateMsg constructCalculatedFieldUpdatedMsg(UpdateMsgType msgType, CalculatedField calculatedField) { public static CalculatedFieldUpdateMsg constructCalculatedFieldUpdatedMsg(UpdateMsgType msgType, CalculatedField calculatedField) {
resetVersion(calculatedField);
return CalculatedFieldUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(calculatedField)) return CalculatedFieldUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(calculatedField))
.setIdMSB(calculatedField.getId().getId().getMostSignificantBits()) .setIdMSB(calculatedField.getId().getId().getMostSignificantBits())
.setIdLSB(calculatedField.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(calculatedField.getId().getId().getLeastSignificantBits()).build();
@ -707,6 +728,7 @@ public class EdgeMsgConstructorUtils {
} }
public static AiModelUpdateMsg constructAiModelUpdatedMsg(UpdateMsgType msgType, AiModel aiModel) { public static AiModelUpdateMsg constructAiModelUpdatedMsg(UpdateMsgType msgType, AiModel aiModel) {
resetVersion(aiModel);
return AiModelUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(aiModel)) return AiModelUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(aiModel))
.setIdMSB(aiModel.getId().getId().getMostSignificantBits()) .setIdMSB(aiModel.getId().getId().getMostSignificantBits())
.setIdLSB(aiModel.getId().getId().getLeastSignificantBits()).build(); .setIdLSB(aiModel.getId().getId().getLeastSignificantBits()).build();

43
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<Void> {
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);
}
}

37
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -135,9 +135,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Lazy @Lazy
private EdgeContextComponent ctx; private EdgeContextComponent ctx;
@Autowired
private TelemetrySubscriptionService tsSubService;
@Autowired @Autowired
private TbClusterService clusterService; 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) { private void save(TenantId tenantId, EdgeId edgeId, String key, long value) {
log.debug("[{}][{}] Updating long edge telemetry [{}] [{}]", tenantId, edgeId, key, value); log.debug("[{}][{}] Updating long edge telemetry [{}] [{}]", tenantId, edgeId, key, value);
if (persistToTelemetry) { if (persistToTelemetry) {
tsSubService.saveTimeseries(TimeseriesSaveRequest.builder() ctx.getTsSubService().saveTimeseries(TimeseriesSaveRequest.builder()
.tenantId(tenantId) .tenantId(tenantId)
.entityId(edgeId) .entityId(edgeId)
.entry(new LongDataEntry(key, value)) .entry(new LongDataEntry(key, value))
.callback(new AttributeSaveCallback(tenantId, edgeId, key, value)) .callback(new AttributeSaveCallback(tenantId, edgeId, key, value))
.build()); .build());
} else { } else {
tsSubService.saveAttributes(AttributesSaveRequest.builder() ctx.getTsSubService().saveAttributes(AttributesSaveRequest.builder()
.tenantId(tenantId) .tenantId(tenantId)
.entityId(edgeId) .entityId(edgeId)
.scope(AttributeScope.SERVER_SCOPE) .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) { private void save(TenantId tenantId, EdgeId edgeId, String key, boolean value) {
log.debug("[{}][{}] Updating boolean edge telemetry [{}] [{}]", tenantId, edgeId, key, value); log.debug("[{}][{}] Updating boolean edge telemetry [{}] [{}]", tenantId, edgeId, key, value);
if (persistToTelemetry) { if (persistToTelemetry) {
tsSubService.saveTimeseries(TimeseriesSaveRequest.builder() ctx.getTsSubService().saveTimeseries(TimeseriesSaveRequest.builder()
.tenantId(tenantId) .tenantId(tenantId)
.entityId(edgeId) .entityId(edgeId)
.entry(new BooleanDataEntry(key, value)) .entry(new BooleanDataEntry(key, value))
.callback(new AttributeSaveCallback(tenantId, edgeId, key, value)) .callback(new AttributeSaveCallback(tenantId, edgeId, key, value))
.build()); .build());
} else { } else {
tsSubService.saveAttributes(AttributesSaveRequest.builder() ctx.getTsSubService().saveAttributes(AttributesSaveRequest.builder()
.tenantId(tenantId) .tenantId(tenantId)
.entityId(edgeId) .entityId(edgeId)
.scope(AttributeScope.SERVER_SCOPE) .scope(AttributeScope.SERVER_SCOPE)
@ -590,32 +587,6 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} }
} }
private static class AttributeSaveCallback implements FutureCallback<Void> {
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) { private void pushRuleEngineMessage(TenantId tenantId, Edge edge, long ts, TbMsgType msgType) {
try { try {
EdgeId edgeId = edge.getId(); EdgeId edgeId = edge.getId();

25
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 lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable; import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.data.util.Pair; 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.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EdgeUtils; 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.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.AttributesSaveResult; import org.thingsboard.server.common.data.kv.AttributesSaveResult;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; 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.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.limit.LimitedApi; import org.thingsboard.server.common.data.limit.LimitedApi;
@ -195,7 +197,7 @@ public abstract class EdgeGrpcSession implements Closeable {
} }
startSyncProcess(fullSync); startSyncProcess(fullSync);
} else { } else {
syncInProgress = false; updateSyncInProgress(false);
} }
} }
if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) { if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) {
@ -241,6 +243,11 @@ public abstract class EdgeGrpcSession implements Closeable {
log.debug("[{}] onConfigurationUpdate [{}]", sessionId, edge); log.debug("[{}] onConfigurationUpdate [{}]", sessionId, edge);
this.tenantId = edge.getTenantId(); this.tenantId = edge.getTenantId();
this.edge = edge; 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() EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder()
.setConfiguration(EdgeMsgConstructorUtils.constructEdgeConfiguration(edge)).build(); .setConfiguration(EdgeMsgConstructorUtils.constructEdgeConfiguration(edge)).build();
ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder() ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder()
@ -252,7 +259,7 @@ public abstract class EdgeGrpcSession implements Closeable {
public void startSyncProcess(boolean fullSync) { public void startSyncProcess(boolean fullSync) {
if (!syncInProgress) { if (!syncInProgress) {
log.info("[{}][{}][{}] Staring edge sync process", tenantId, edge.getId(), sessionId); log.info("[{}][{}][{}] Staring edge sync process", tenantId, edge.getId(), sessionId);
syncInProgress = true; updateSyncInProgress(true);
interruptGeneralProcessingOnSync(); interruptGeneralProcessingOnSync();
doSync(new EdgeSyncCursor(ctx, edge, fullSync)); doSync(new EdgeSyncCursor(ctx, edge, fullSync));
} else { } else {
@ -398,6 +405,18 @@ public abstract class EdgeGrpcSession implements Closeable {
ctx.getAttributesService().save(tenantId, edge.getId(), AttributeScope.SERVER_SCOPE, attributeKvEntry); 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() { private void interruptGeneralProcessingOnSync() {
log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", tenantId, edge.getId(), sessionId); log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", tenantId, edge.getId(), sessionId);
stopCurrentSendDownlinkMsgsTask(true); stopCurrentSendDownlinkMsgsTask(true);
@ -766,7 +785,7 @@ public abstract class EdgeGrpcSession implements Closeable {
} }
private void markSyncCompletedSendEdgeEventUpdate() { private void markSyncCompletedSendEdgeEventUpdate() {
syncInProgress = false; updateSyncInProgress(false);
ctx.getClusterService().onEdgeEventUpdate(new EdgeEventUpdateMsg(edge.getTenantId(), edge.getId())); ctx.getClusterService().onEdgeEventUpdate(new EdgeEventUpdateMsg(edge.getTenantId(), edge.getId()));
} }

4
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.Customer;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EntityId; 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.EdgeContextComponent;
import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher; 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.AiModelEdgeEventFetcher;
@ -63,7 +64,8 @@ public class EdgeSyncCursor {
fetchers.add(new TenantEdgeEventFetcher(ctx.getTenantService())); fetchers.add(new TenantEdgeEventFetcher(ctx.getTenantService()));
fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService()));
fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService())); 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 TenantAdminUsersEdgeEventFetcher(ctx.getUserService()));
fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService())); fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService()));
fetchers.add(new SystemWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService())); fetchers.add(new SystemWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService()));

34
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.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType; 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.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.settings.AdminSettingsService; import org.thingsboard.server.dao.settings.AdminSettingsService;
import java.util.ArrayList;
import java.util.List;
@AllArgsConstructor @AllArgsConstructor
@Slf4j @Slf4j
public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { public class AdminSettingsEdgeEventFetcher extends BasePageableEdgeEventFetcher<AdminSettings> {
private final AdminSettingsService adminSettingsService; private final AdminSettingsService adminSettingsService;
private final TenantId fetcherTenantId;
@Override @Override
public PageLink getPageLink(int pageSize) { PageData<AdminSettings> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) {
return null; return adminSettingsService.findAllByTenantId(fetcherTenantId, pageLink);
}
public PageData<EdgeEvent> fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) {
List<EdgeEvent> 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);
} }
private List<EdgeEvent> fetchAdminSettingsForKeys(TenantId tenantId, EdgeId edgeId, List<String> keys) { @Override
List<EdgeEvent> result = new ArrayList<>(); EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, AdminSettings adminSettings) {
for (String key : keys) { return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS,
AdminSettings adminSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, key); EdgeEventActionType.UPDATED, adminSettings.getId(), null);
if (adminSettings != null) {
result.add(EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.ADMIN_SETTINGS,
EdgeEventActionType.UPDATED, null, JacksonUtil.valueToTree(adminSettings)));
}
}
return result;
} }
} }

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java

@ -41,7 +41,7 @@ public class OAuth2EdgeEventFetcher extends BasePageableEdgeEventFetcher<DomainI
@Override @Override
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, DomainInfo domainInfo) { EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, DomainInfo domainInfo) {
return EdgeUtils.constructEdgeEvent(TenantId.SYS_TENANT_ID, edge.getId(), EdgeEventType.DOMAIN, return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DOMAIN,
EdgeEventActionType.ADDED, domainInfo.getId(), null); EdgeEventActionType.ADDED, domainInfo.getId(), null);
} }

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java

@ -412,6 +412,9 @@ public abstract class BaseEdgeProcessor implements EdgeProcessor {
} }
protected boolean isSaveRequired(HasVersion current, HasVersion updated) { protected boolean isSaveRequired(HasVersion current, HasVersion updated) {
if (current != null) {
current.setVersion(null);
}
updated.setVersion(null); updated.setVersion(null);
return !updated.equals(current); return !updated.equals(current);
} }

10
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/settings/AdminSettingsEdgeProcessor.java

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.AdminSettingsId;
import org.thingsboard.server.gen.edge.v1.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.v1.AdminSettingsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg;
import org.thingsboard.server.gen.edge.v1.EdgeVersion; import org.thingsboard.server.gen.edge.v1.EdgeVersion;
@ -35,7 +36,14 @@ public class AdminSettingsEdgeProcessor extends BaseEdgeProcessor {
@Override @Override
public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) {
AdminSettings adminSettings = JacksonUtil.convertValue(edgeEvent.getBody(), AdminSettings.class); AdminSettings adminSettings = null;
if (edgeEvent.getEntityId() != null) {
AdminSettingsId adminSettingsId = new AdminSettingsId(edgeEvent.getEntityId());
adminSettings = edgeCtx.getAdminSettingsService().findAdminSettingsById(edgeEvent.getTenantId(), adminSettingsId);
} else if (edgeEvent.getBody() != null && !edgeEvent.getBody().isEmpty()) {
// legacy
adminSettings = JacksonUtil.convertValue(edgeEvent.getBody(), AdminSettings.class);
}
if (adminSettings == null) { if (adminSettings == null) {
return null; return null;
} }

4
application/src/main/resources/thingsboard.yml

@ -78,6 +78,10 @@ server:
max_entities_per_data_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_DATA_SUBSCRIPTION:10000}" max_entities_per_data_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_DATA_SUBSCRIPTION:10000}"
# Maximum number of alarms returned for single alarm subscription. For example, no more than 10,000 alarms on the alarm widget # Maximum number of alarms returned for single alarm subscription. For example, no more than 10,000 alarms on the alarm widget
max_entities_per_alarm_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_ALARM_SUBSCRIPTION:10000}" max_entities_per_alarm_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_ALARM_SUBSCRIPTION:10000}"
# Maximum size (bytes) of incoming WS text message
max_text_message_buffer_size: "${TB_SERVER_WS_MAX_TEXT_MESSAGE_BUFFER_SIZE:32768}"
# Maximum size (bytes) of incoming WS binary message
max_binary_message_buffer_size: "${TB_SERVER_WS_MAX_BINARY_MESSAGE_BUFFER_SIZE:32768}"
# Maximum queue size of the websocket updates per session. This restriction prevents infinite updates of WS # Maximum queue size of the websocket updates per session. This restriction prevents infinite updates of WS
max_queue_messages_per_session: "${TB_SERVER_WS_DEFAULT_QUEUE_MESSAGES_PER_SESSION:1000}" max_queue_messages_per_session: "${TB_SERVER_WS_DEFAULT_QUEUE_MESSAGES_PER_SESSION:1000}"
# Maximum time between WS session opening and sending auth command # Maximum time between WS session opening and sending auth command

9
application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java

@ -39,6 +39,7 @@ import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceInfo; import org.thingsboard.server.common.data.DeviceInfo;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.HasVersion;
import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.SaveOtaPackageInfoRequest; import org.thingsboard.server.common.data.SaveOtaPackageInfoRequest;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
@ -605,7 +606,7 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true); DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true);
Assert.assertNotNull(deviceCredentialsMsg); Assert.assertNotNull(deviceCredentialsMsg);
Assert.assertEquals(savedDevice.getId(), deviceCredentialsMsg.getDeviceId()); Assert.assertEquals(savedDevice.getId(), deviceCredentialsMsg.getDeviceId());
Assert.assertEquals(deviceCredentials, deviceCredentialsMsg); compareHasVersionEntities(deviceCredentials, deviceCredentialsMsg);
return savedDevice; return savedDevice;
} }
@ -775,4 +776,10 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
}); });
} }
protected void compareHasVersionEntities(HasVersion entity1, HasVersion entity2) {
entity1.setVersion(null);
entity2.setVersion(null);
Assert.assertEquals(entity1, entity2);
}
} }

79
application/src/test/java/org/thingsboard/server/edge/AdminSettingsEdgeTest.java

@ -0,0 +1,79 @@
/**
* 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.edge;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.protobuf.AbstractMessage;
import org.junit.Assert;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.gen.edge.v1.AdminSettingsUpdateMsg;
@DaoSqlTest
public class AdminSettingsEdgeTest extends AbstractEdgeTest {
@Autowired
private AdminSettingsService adminSettingsService;
@Test
public void testAdminSettings() throws Exception {
loginSysAdmin();
// save
AdminSettings adminSettings = new AdminSettings();
adminSettings.setKey("edgeTest");
ObjectNode jsonValue = JacksonUtil.newObjectNode();
jsonValue.put("key1", "value1");
adminSettings.setJsonValue(jsonValue);
edgeImitator.expectMessageAmount(1);
AdminSettings savedAdminSettings = doPost("/api/admin/settings", adminSettings, AdminSettings.class);
Assert.assertTrue(edgeImitator.waitForMessages());
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof AdminSettingsUpdateMsg);
AdminSettingsUpdateMsg adminSettingsUpdateMsg = (AdminSettingsUpdateMsg) latestMessage;
AdminSettings adminSettingsMsg = JacksonUtil.fromString(adminSettingsUpdateMsg.getEntity(), AdminSettings.class, true);
Assert.assertNotNull(adminSettingsMsg);
Assert.assertEquals("edgeTest", adminSettingsMsg.getKey());
Assert.assertEquals("value1", adminSettingsMsg.getJsonValue().get("key1").asText());
// update
ObjectNode updatedJsonValue = (ObjectNode) savedAdminSettings.getJsonValue();
updatedJsonValue.put("key2", "value2");
savedAdminSettings.setJsonValue(updatedJsonValue);
edgeImitator.expectMessageAmount(1);
doPost("/api/admin/settings", savedAdminSettings, AdminSettings.class);
Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof AdminSettingsUpdateMsg);
adminSettingsUpdateMsg = (AdminSettingsUpdateMsg) latestMessage;
adminSettingsMsg = JacksonUtil.fromString(adminSettingsUpdateMsg.getEntity(), AdminSettings.class, true);
Assert.assertNotNull(adminSettingsMsg);
Assert.assertEquals("edgeTest", adminSettingsMsg.getKey());
Assert.assertEquals("value1", adminSettingsMsg.getJsonValue().get("key1").asText());
Assert.assertEquals("value2", adminSettingsMsg.getJsonValue().get("key2").asText());
adminSettingsService.deleteAdminSettingsByTenantIdAndKey(savedAdminSettings.getTenantId(), "edgeTest");
}
}

4
application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java

@ -58,7 +58,7 @@ public class AssetEdgeTest extends AbstractEdgeTest {
Asset assetMsg = JacksonUtil.fromString(assetUpdateMsg.getEntity(), Asset.class, true); Asset assetMsg = JacksonUtil.fromString(assetUpdateMsg.getEntity(), Asset.class, true);
Assert.assertNotNull(assetMsg); Assert.assertNotNull(assetMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
Assert.assertEquals(savedAsset, assetMsg); compareHasVersionEntities(savedAsset, assetMsg);
Optional<AssetProfileUpdateMsg> assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class); Optional<AssetProfileUpdateMsg> assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class);
Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent()); Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent());
AssetProfileUpdateMsg assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get(); AssetProfileUpdateMsg assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get();
@ -109,7 +109,7 @@ public class AssetEdgeTest extends AbstractEdgeTest {
assetMsg = JacksonUtil.fromString(assetUpdateMsg.getEntity(), Asset.class, true); assetMsg = JacksonUtil.fromString(assetUpdateMsg.getEntity(), Asset.class, true);
Assert.assertNotNull(assetMsg); Assert.assertNotNull(assetMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
Assert.assertEquals(savedAsset, assetMsg); compareHasVersionEntities(savedAsset, assetMsg);
assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class); assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class);
Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent()); Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent());
assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get(); assetProfileUpdateMsg = assetProfileUpdateMsgOpt.get();

2
application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java

@ -53,7 +53,7 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest {
AssetProfileUpdateMsg assetProfileUpdateMsg = (AssetProfileUpdateMsg) latestMessage; AssetProfileUpdateMsg assetProfileUpdateMsg = (AssetProfileUpdateMsg) latestMessage;
AssetProfile assetProfileMsg = JacksonUtil.fromString(assetProfileUpdateMsg.getEntity(), AssetProfile.class, true); AssetProfile assetProfileMsg = JacksonUtil.fromString(assetProfileUpdateMsg.getEntity(), AssetProfile.class, true);
Assert.assertNotNull(assetProfileMsg); Assert.assertNotNull(assetProfileMsg);
Assert.assertEquals(assetProfile, assetProfileMsg); compareHasVersionEntities(assetProfile, assetProfileMsg);
Assert.assertEquals("Building", assetProfileMsg.getName()); Assert.assertEquals("Building", assetProfileMsg.getName());
Assert.assertEquals(buildingsRuleChainId, assetProfileMsg.getDefaultEdgeRuleChainId()); Assert.assertEquals(buildingsRuleChainId, assetProfileMsg.getDefaultEdgeRuleChainId());
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetProfileUpdateMsg.getMsgType());

2
application/src/test/java/org/thingsboard/server/edge/CalculatedFieldEdgeTest.java

@ -147,7 +147,7 @@ public class CalculatedFieldEdgeTest extends AbstractEdgeTest {
CalculatedFieldUpdateMsg calculatedFieldUpdateMsg = (CalculatedFieldUpdateMsg) latestMessage; CalculatedFieldUpdateMsg calculatedFieldUpdateMsg = (CalculatedFieldUpdateMsg) latestMessage;
CalculatedField calculatedFieldFromEdge = JacksonUtil.fromString(calculatedFieldUpdateMsg.getEntity(), CalculatedField.class, true); CalculatedField calculatedFieldFromEdge = JacksonUtil.fromString(calculatedFieldUpdateMsg.getEntity(), CalculatedField.class, true);
Assert.assertNotNull(calculatedFieldFromEdge); Assert.assertNotNull(calculatedFieldFromEdge);
Assert.assertEquals(savedCalculatedField, calculatedFieldFromEdge); compareHasVersionEntities(savedCalculatedField, calculatedFieldFromEdge);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, calculatedFieldUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, calculatedFieldUpdateMsg.getMsgType());
} }

4
application/src/test/java/org/thingsboard/server/edge/CustomerEdgeTest.java

@ -60,7 +60,7 @@ public class CustomerEdgeTest extends AbstractEdgeTest {
CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get();
Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType());
Assert.assertEquals(savedCustomer, customerMsg); compareHasVersionEntities(savedCustomer, customerMsg);
testAutoGeneratedCodeByProtobuf(customerUpdateMsg); testAutoGeneratedCodeByProtobuf(customerUpdateMsg);
// update customer // update customer
@ -73,7 +73,7 @@ public class CustomerEdgeTest extends AbstractEdgeTest {
customerUpdateMsg = (CustomerUpdateMsg) latestMessage; customerUpdateMsg = (CustomerUpdateMsg) latestMessage;
customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true);
Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, customerUpdateMsg.getMsgType());
Assert.assertEquals(savedCustomer, customerMsg); compareHasVersionEntities(savedCustomer, customerMsg);
// delete customer // delete customer
edgeImitator.expectMessageAmount(2); edgeImitator.expectMessageAmount(2);

2
application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java

@ -201,7 +201,7 @@ public class DashboardEdgeTest extends AbstractEdgeTest {
CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get();
Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType());
Assert.assertEquals(savedCustomer, customerMsg); compareHasVersionEntities(savedCustomer, customerMsg);
Dashboard dashboard = buildDashboardForUplinkMsg(savedCustomer); Dashboard dashboard = buildDashboardForUplinkMsg(savedCustomer);

8
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); Device deviceFromMsg = JacksonUtil.fromString(deviceUpdateMsg.getEntity(), Device.class, true);
Assert.assertNotNull(deviceFromMsg); Assert.assertNotNull(deviceFromMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); 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.getId(), deviceFromMsg.getId());
Assert.assertEquals(savedDevice.getName(), deviceFromMsg.getName()); Assert.assertEquals(savedDevice.getName(), deviceFromMsg.getName());
Assert.assertEquals(savedDevice.getType(), deviceFromMsg.getType()); Assert.assertEquals(savedDevice.getType(), deviceFromMsg.getType());
@ -222,7 +222,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest {
Assert.assertTrue(latestMessage instanceof DeviceCredentialsUpdateMsg); Assert.assertTrue(latestMessage instanceof DeviceCredentialsUpdateMsg);
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = (DeviceCredentialsUpdateMsg) latestMessage; DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = (DeviceCredentialsUpdateMsg) latestMessage;
DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true); DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true);
Assert.assertEquals(deviceCredentials, deviceCredentialsMsg); compareHasVersionEntities(deviceCredentials, deviceCredentialsMsg);
// update device credentials - X509_CERTIFICATE // update device credentials - X509_CERTIFICATE
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
@ -272,7 +272,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest {
Device deviceMsg = JacksonUtil.fromString(deviceUpdateMsg.getEntity(), Device.class, true); Device deviceMsg = JacksonUtil.fromString(deviceUpdateMsg.getEntity(), Device.class, true);
Assert.assertNotNull(deviceMsg); Assert.assertNotNull(deviceMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
Assert.assertEquals(savedDevice, deviceMsg); compareHasVersionEntities(savedDevice, deviceMsg);
Assert.assertEquals(firmwareOtaPackageInfo.getId(), deviceMsg.getFirmwareId()); Assert.assertEquals(firmwareOtaPackageInfo.getId(), deviceMsg.getFirmwareId());
Assert.assertEquals(softwareOtaPackageInfo.getId(), deviceMsg.getSoftwareId()); Assert.assertEquals(softwareOtaPackageInfo.getId(), deviceMsg.getSoftwareId());
deviceData = deviceMsg.getDeviceData(); deviceData = deviceMsg.getDeviceData();
@ -387,7 +387,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest {
DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true); DeviceCredentials deviceCredentialsMsg = JacksonUtil.fromString(deviceCredentialsUpdateMsg.getEntity(), DeviceCredentials.class, true);
Assert.assertNotNull(deviceCredentialsMsg); Assert.assertNotNull(deviceCredentialsMsg);
Assert.assertEquals(device.getId(), deviceCredentialsMsg.getDeviceId()); Assert.assertEquals(device.getId(), deviceCredentialsMsg.getDeviceId());
Assert.assertEquals(deviceCredentials, deviceCredentialsMsg); compareHasVersionEntities(deviceCredentials, deviceCredentialsMsg);
} }
@Test @Test

12
application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java

@ -83,7 +83,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
// update device profile // update device profile
@ -108,7 +108,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
// delete profile // delete profile
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
@ -146,7 +146,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
// delete profile when edge is offline // delete profile when edge is offline
@ -186,7 +186,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
Assert.assertEquals(DeviceTransportType.SNMP, deviceProfileMsg.getTransportType()); Assert.assertEquals(DeviceTransportType.SNMP, deviceProfileMsg.getTransportType());
@ -224,7 +224,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
Assert.assertEquals(DeviceTransportType.LWM2M, deviceProfileMsg.getTransportType()); Assert.assertEquals(DeviceTransportType.LWM2M, deviceProfileMsg.getTransportType());
@ -273,7 +273,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true); DeviceProfile deviceProfileMsg = JacksonUtil.fromString(deviceProfileUpdateMsg.getEntity(), DeviceProfile.class, true);
Assert.assertNotNull(deviceProfileMsg); Assert.assertNotNull(deviceProfileMsg);
Assert.assertEquals(deviceProfile, deviceProfileMsg); compareHasVersionEntities(deviceProfile, deviceProfileMsg);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType());
Assert.assertEquals(DeviceTransportType.COAP, deviceProfileMsg.getTransportType()); Assert.assertEquals(DeviceTransportType.COAP, deviceProfileMsg.getTransportType());

27
application/src/test/java/org/thingsboard/server/edge/EdgeTest.java

@ -15,10 +15,12 @@
*/ */
package org.thingsboard.server.edge; package org.thingsboard.server.edge;
import com.fasterxml.jackson.databind.JsonNode;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Customer; 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.edge.Edge;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId; 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.CustomerUpdateMsg;
import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; import org.thingsboard.server.gen.edge.v1.EdgeConfiguration;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; 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.Optional;
import java.util.UUID; import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
@DaoSqlTest @DaoSqlTest
public class EdgeTest extends AbstractEdgeTest { public class EdgeTest extends AbstractEdgeTest {
@ -57,7 +63,7 @@ public class EdgeTest extends AbstractEdgeTest {
CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get(); CustomerUpdateMsg customerUpdateMsg = customerUpdateOpt.get();
Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true); Customer customerMsg = JacksonUtil.fromString(customerUpdateMsg.getEntity(), Customer.class, true);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, customerUpdateMsg.getMsgType());
Assert.assertEquals(savedCustomer, customerMsg); compareHasVersionEntities(savedCustomer, customerMsg);
// unassign edge from customer // unassign edge from customer
edgeImitator.expectMessageAmount(2); edgeImitator.expectMessageAmount(2);
@ -76,4 +82,23 @@ public class EdgeTest extends AbstractEdgeTest {
Assert.assertEquals(savedCustomer.getUuidId().getMostSignificantBits(), customerUpdateMsg.getIdMSB()); Assert.assertEquals(savedCustomer.getUuidId().getMostSignificantBits(), customerUpdateMsg.getIdMSB());
Assert.assertEquals(savedCustomer.getUuidId().getLeastSignificantBits(), customerUpdateMsg.getIdLSB()); 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);
}
} }

4
application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java

@ -68,7 +68,7 @@ public class EntityViewEdgeTest extends AbstractEdgeTest {
EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage; EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage;
EntityView entityView = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true); EntityView entityView = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true);
Assert.assertNotNull(entityView); Assert.assertNotNull(entityView);
Assert.assertEquals(savedEntityView, entityView); compareHasVersionEntities(savedEntityView, entityView);
Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType());
// request entity view(s) for device // request entity view(s) for device
@ -265,7 +265,7 @@ public class EntityViewEdgeTest extends AbstractEdgeTest {
EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage; EntityViewUpdateMsg entityViewUpdateMsg = (EntityViewUpdateMsg) latestMessage;
EntityView entityViewMsg = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true); EntityView entityViewMsg = JacksonUtil.fromString(entityViewUpdateMsg.getEntity(), EntityView.class, true);
Assert.assertNotNull(entityViewMsg); Assert.assertNotNull(entityViewMsg);
Assert.assertEquals(entityView, entityViewMsg); compareHasVersionEntities(entityView, entityViewMsg);
Assert.assertEquals(device.getId(), entityViewMsg.getEntityId()); Assert.assertEquals(device.getId(), entityViewMsg.getEntityId());
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, entityViewUpdateMsg.getMsgType());
testAutoGeneratedCodeByProtobuf(entityViewUpdateMsg); testAutoGeneratedCodeByProtobuf(entityViewUpdateMsg);

10
application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java

@ -57,7 +57,7 @@ public class RelationEdgeTest extends AbstractEdgeTest {
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true);
Assert.assertNotNull(entityRelation); Assert.assertNotNull(entityRelation);
Assert.assertEquals(relation, entityRelation); compareHasVersionEntities(relation, entityRelation);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
// delete relation // delete relation
@ -76,7 +76,7 @@ public class RelationEdgeTest extends AbstractEdgeTest {
relationUpdateMsg = (RelationUpdateMsg) latestMessage; relationUpdateMsg = (RelationUpdateMsg) latestMessage;
entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true);
Assert.assertNotNull(entityRelation); Assert.assertNotNull(entityRelation);
Assert.assertEquals(deletedRelation, entityRelation); compareHasVersionEntities(deletedRelation, entityRelation);
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
} }
@ -155,7 +155,7 @@ public class RelationEdgeTest extends AbstractEdgeTest {
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true);
Assert.assertNotNull(entityRelation); Assert.assertNotNull(entityRelation);
Assert.assertEquals(deviceToAssetRelation, entityRelation); compareHasVersionEntities(deviceToAssetRelation, entityRelation);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
} }
@ -177,7 +177,7 @@ public class RelationEdgeTest extends AbstractEdgeTest {
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); EntityRelation entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true);
Assert.assertNotNull(entityRelation); Assert.assertNotNull(entityRelation);
Assert.assertEquals(relation, entityRelation); compareHasVersionEntities(relation, entityRelation);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
// delete relation // delete relation
@ -196,7 +196,7 @@ public class RelationEdgeTest extends AbstractEdgeTest {
relationUpdateMsg = (RelationUpdateMsg) latestMessage; relationUpdateMsg = (RelationUpdateMsg) latestMessage;
entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true); entityRelation = JacksonUtil.fromString(relationUpdateMsg.getEntity(), EntityRelation.class, true);
Assert.assertNotNull(entityRelation); Assert.assertNotNull(entityRelation);
Assert.assertEquals(deletedRelation, entityRelation); compareHasVersionEntities(deletedRelation, entityRelation);
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
} }

4
application/src/test/java/org/thingsboard/server/edge/TenantEdgeTest.java

@ -50,7 +50,7 @@ public class TenantEdgeTest extends AbstractEdgeTest {
TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileUpdateMsgOpt.get(); TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileUpdateMsgOpt.get();
Tenant tenantMsg = JacksonUtil.fromString(tenantUpdateMsg.getEntity(), Tenant.class, true); Tenant tenantMsg = JacksonUtil.fromString(tenantUpdateMsg.getEntity(), Tenant.class, true);
Assert.assertNotNull(tenantMsg); Assert.assertNotNull(tenantMsg);
Assert.assertEquals(savedTenant, tenantMsg); compareHasVersionEntities(savedTenant, tenantMsg);
TenantProfile tenantProfileMsg = JacksonUtil.fromString(tenantProfileUpdateMsg.getEntity(), TenantProfile.class, true); TenantProfile tenantProfileMsg = JacksonUtil.fromString(tenantProfileUpdateMsg.getEntity(), TenantProfile.class, true);
Assert.assertNotNull(tenantProfileMsg); Assert.assertNotNull(tenantProfileMsg);
Assert.assertEquals(tenantMsg.getTenantProfileId(), tenantProfileMsg.getId()); Assert.assertEquals(tenantMsg.getTenantProfileId(), tenantProfileMsg.getId());
@ -75,7 +75,7 @@ public class TenantEdgeTest extends AbstractEdgeTest {
Assert.assertNotNull(tenantProfileMsg); Assert.assertNotNull(tenantProfileMsg);
// tenant update // tenant update
Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, tenantUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, tenantUpdateMsg.getMsgType());
Assert.assertEquals(savedTenant, tenantMsg); compareHasVersionEntities(savedTenant, tenantMsg);
Assert.assertEquals(savedTenant.getTenantProfileId(), tenantProfileMsg.getId()); Assert.assertEquals(savedTenant.getTenantProfileId(), tenantProfileMsg.getId());
} }

2
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); WidgetType widgetsType = JacksonUtil.fromString(widgetTypeUpdateMsg.getEntity(), WidgetType.class, true);
Assert.assertNotNull(widgetsType); Assert.assertNotNull(widgetsType);
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, widgetTypeUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, widgetTypeUpdateMsg.getMsgType());
Assert.assertEquals(savedWidgetType, widgetsType); compareHasVersionEntities(savedWidgetType, widgetsType);
// update widget bundle // update widget bundle
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);

181
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.TbCalculatedFieldsNode;
import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode;
import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; 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.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType; 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.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.RuleChainMetaData;
import org.thingsboard.server.common.data.rule.RuleNode; 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.EdgeVersion;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
@ -278,4 +296,167 @@ public class EdgeMsgConstructorUtilsTest {
edgeEvent.setBody(body); edgeEvent.setBody(body);
return edgeEvent; 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");
}
} }

4
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.AdminSettings;
import org.thingsboard.server.common.data.id.AdminSettingsId; import org.thingsboard.server.common.data.id.AdminSettingsId;
import org.thingsboard.server.common.data.id.TenantId; 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; import org.thingsboard.server.dao.entity.EntityDaoService;
public interface AdminSettingsService extends EntityDaoService { public interface AdminSettingsService extends EntityDaoService {
@ -28,6 +30,8 @@ public interface AdminSettingsService extends EntityDaoService {
AdminSettings findAdminSettingsByTenantIdAndKey(TenantId tenantId, String key); AdminSettings findAdminSettingsByTenantIdAndKey(TenantId tenantId, String key);
PageData<AdminSettings> findAllByTenantId(TenantId tenantId, PageLink pageLink);
AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings); AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings);
boolean deleteAdminSettingsByTenantIdAndKey(TenantId tenantId, String key); boolean deleteAdminSettingsByTenantIdAndKey(TenantId tenantId, String key);

1
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 EDGE_MSG_SOURCE = "edge";
public static final String MSG_SOURCE_KEY = "source"; public static final String MSG_SOURCE_KEY = "source";
public static final String EDGE_VERSION_ATTR_KEY = "edgeVersion"; 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"; public static final String LAST_CONNECTED_GATEWAY = "lastConnectedGateway";

2
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), TENANT_PROFILE(true, EntityType.TENANT_PROFILE),
WIDGETS_BUNDLE(true, EntityType.WIDGETS_BUNDLE), WIDGETS_BUNDLE(true, EntityType.WIDGETS_BUNDLE),
WIDGET_TYPE(true, EntityType.WIDGET_TYPE), WIDGET_TYPE(true, EntityType.WIDGET_TYPE),
ADMIN_SETTINGS(true, null), ADMIN_SETTINGS(true, EntityType.ADMIN_SETTINGS),
OTA_PACKAGE(true, EntityType.OTA_PACKAGE), OTA_PACKAGE(true, EntityType.OTA_PACKAGE),
QUEUE(true, EntityType.QUEUE), QUEUE(true, EntityType.QUEUE),
NOTIFICATION_RULE(true, EntityType.NOTIFICATION_RULE), NOTIFICATION_RULE(true, EntityType.NOTIFICATION_RULE),

1
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 DOMAIN -> new DomainId(uuid);
case CALCULATED_FIELD -> new CalculatedFieldId(uuid); case CALCULATED_FIELD -> new CalculatedFieldId(uuid);
case AI_MODEL -> new AiModelId(uuid); case AI_MODEL -> new AiModelId(uuid);
case ADMIN_SETTINGS -> new AdminSettingsId(uuid);
default -> throw new IllegalArgumentException("EdgeEventType " + edgeEventType + " is not supported!"); default -> throw new IllegalArgumentException("EdgeEventType " + edgeEventType + " is not supported!");
}; };
} }

13
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.assertThatNoException;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
@ResourceLock("SsrfProtectionValidatorTest") // some tests mutate static additional-blocked-hosts
public class SsrfProtectionValidatorTest { 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 @ParameterizedTest
@ValueSource(strings = { @ValueSource(strings = {
"http://example.com", "http://example.com",
@ -207,7 +204,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedSingleIp() { void testAdditionalBlockedSingleIp() {
try { try {
SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("8.8.8.8")); SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("8.8.8.8"));
@ -222,7 +218,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedCidrSlash10() { void testAdditionalBlockedCidrSlash10() {
try { try {
// Use 44.0.0.0/10 (not blocked by default) to verify CIDR /10 matching // Use 44.0.0.0/10 (not blocked by default) to verify CIDR /10 matching
@ -243,7 +238,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedCidrSlash24() { void testAdditionalBlockedCidrSlash24() {
try { try {
SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("198.51.100.0/24")); SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("198.51.100.0/24"));
@ -261,7 +255,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedHostnameViaValidateUri() { void testAdditionalBlockedHostnameViaValidateUri() {
try { try {
SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("evil.corp")); SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("evil.corp"));
@ -275,7 +268,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedHostnameCaseInsensitive() { void testAdditionalBlockedHostnameCaseInsensitive() {
try { try {
SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("My-Service.Corp")); SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("My-Service.Corp"));
@ -289,7 +281,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testSetAdditionalBlockedHostsEmptyAndNull() { void testSetAdditionalBlockedHostsEmptyAndNull() {
// Should not throw // Should not throw
SsrfProtectionValidator.setAdditionalBlockedHosts(Collections.emptyList()); SsrfProtectionValidator.setAdditionalBlockedHosts(Collections.emptyList());
@ -307,7 +298,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedCidrViaValidateUri() { void testAdditionalBlockedCidrViaValidateUri() {
// 203.0.113.0/24 (TEST-NET-3) is not blocked by default // 203.0.113.0/24 (TEST-NET-3) is not blocked by default
URI uri = URI.create("http://203.0.113.1"); URI uri = URI.create("http://203.0.113.1");
@ -323,7 +313,6 @@ public class SsrfProtectionValidatorTest {
} }
@Test @Test
@ResourceLock(SYNC_LOCK)
void testAdditionalBlockedMixedConfig() { void testAdditionalBlockedMixedConfig() {
try { try {
SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("203.0.113.0/24", "evil.corp", "8.8.8.8")); SsrfProtectionValidator.setAdditionalBlockedHosts(List.of("203.0.113.0/24", "evil.corp", "8.8.8.8"));

5
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.ActionEntityEvent;
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent;
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; 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.relation.RelationService;
import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.service.DataValidator; 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.sql.JpaExecutorService;
import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.exception.DataValidationException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
@ -239,9 +239,10 @@ public class EdgeServiceImpl extends AbstractCachedEntityService<EdgeCacheKey, E
return edge; return edge;
} }
edge.setCustomerId(customerId); edge.setCustomerId(customerId);
Edge result = saveEdge(edge);
eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).entityId(edgeId) eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).entityId(edgeId)
.body(JacksonUtil.toString(customerId)).actionType(ActionType.ASSIGNED_TO_CUSTOMER).build()); .body(JacksonUtil.toString(customerId)).actionType(ActionType.ASSIGNED_TO_CUSTOMER).build());
return saveEdge(edge); return result;
} }
@Override @Override

4
dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsDao.java

@ -17,6 +17,8 @@ package org.thingsboard.server.dao.settings;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.id.TenantId; 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.Dao; import org.thingsboard.server.dao.Dao;
import java.util.UUID; import java.util.UUID;
@ -27,6 +29,8 @@ public interface AdminSettingsDao extends Dao<AdminSettings> {
AdminSettings findByTenantIdAndKey(UUID tenantId, String key); AdminSettings findByTenantIdAndKey(UUID tenantId, String key);
PageData<AdminSettings> findAllByTenantId(TenantId tenantId, PageLink pageLink);
boolean removeByTenantIdAndKey(UUID tenantId, String key); boolean removeByTenantIdAndKey(UUID tenantId, String key);
void removeByTenantId(UUID tenantId); void removeByTenantId(UUID tenantId);

19
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 com.google.common.util.concurrent.FluentFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.EntityType; 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.EntityId;
import org.thingsboard.server.common.data.id.HasId; import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.TenantId; 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.DataValidator;
import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.service.Validator;
@ -44,6 +48,9 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
@Autowired @Autowired
private DataValidator<AdminSettings> adminSettingsValidator; private DataValidator<AdminSettings> adminSettingsValidator;
@Autowired
protected ApplicationEventPublisher eventPublisher;
@Override @Override
public AdminSettings findAdminSettingsById(TenantId tenantId, AdminSettingsId adminSettingsId) { public AdminSettings findAdminSettingsById(TenantId tenantId, AdminSettingsId adminSettingsId) {
log.trace("Executing findAdminSettingsById [{}]", adminSettingsId); log.trace("Executing findAdminSettingsById [{}]", adminSettingsId);
@ -63,10 +70,15 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
return adminSettingsDao.findByTenantIdAndKey(tenantId.getId(), key); return adminSettingsDao.findByTenantIdAndKey(tenantId.getId(), key);
} }
@Override
public PageData<AdminSettings> findAllByTenantId(TenantId tenantId, PageLink pageLink) {
return adminSettingsDao.findAllByTenantId(tenantId, pageLink);
}
@Override @Override
public AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings) { public AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings) {
log.trace("Executing saveAdminSettings [{}]", adminSettings); log.trace("Executing saveAdminSettings [{}]", adminSettings);
adminSettingsValidator.validate(adminSettings, data -> tenantId); AdminSettings oldAdminSettings = adminSettingsValidator.validate(adminSettings, data -> tenantId);
if (adminSettings.getKey().equals("mail")) { if (adminSettings.getKey().equals("mail")) {
AdminSettings mailSettings = findAdminSettingsByKey(tenantId, "mail"); AdminSettings mailSettings = findAdminSettingsByKey(tenantId, "mail");
if (mailSettings != null) { if (mailSettings != null) {
@ -84,7 +96,10 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
if (adminSettings.getTenantId() == null) { if (adminSettings.getTenantId() == null) {
adminSettings.setTenantId(TenantId.SYS_TENANT_ID); 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 @Override

25
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.common.util.JacksonUtil;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.id.TenantId; 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.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.assertj.core.api.Assertions.assertThatThrownBy;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
@ -113,4 +116,24 @@ public class AdminSettingsServiceTest extends AbstractServiceTest {
}).hasMessageContaining("already exists"); }).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<AdminSettings> pageData = adminSettingsService.findAllByTenantId(tenantId, new PageLink(pageSize));
assertThat(pageData.getData().size()).isEqualTo(pageSize);
assertThat(pageData.getTotalElements()).isEqualTo(totalElements);
}
} }

Loading…
Cancel
Save