Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard into feature/entity-agg-cf

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
6ebdbd6b60
  1. 2
      application/src/main/data/upgrade/basic/schema_update.sql
  2. 17
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  3. 18
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  4. 25
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
  6. 15
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java
  8. 14
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  9. 87
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  10. 13
      application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java
  11. 96
      application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java
  12. 27
      application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java
  13. 2
      application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java
  14. 8
      application/src/test/java/org/thingsboard/server/controller/UserControllerTest.java
  15. 38
      application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
  16. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java
  17. 36
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java
  18. 1
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java
  19. 5
      application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java
  20. 26
      application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/Ota5LwM2MIntegrationTest.java
  21. 7
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test.java
  22. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test.java
  23. 6
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test.java
  24. 54
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java
  25. 5
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationDiscoverTest.java
  26. 5
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationDiscoverWriteAttributesTest.java
  27. 7
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java
  28. 23
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer10Test.java
  29. 20
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer11Test.java
  30. 21
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer12Test.java
  31. 9
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java
  32. 14
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/PskLwm2mIntegrationTest.java
  33. 0
      application/src/test/resources/lwm2m/3-1_2.xml
  34. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java
  35. 4
      common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java
  36. 10
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java
  37. 4
      common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java
  38. 10
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java
  39. 26
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  40. 4
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java
  41. 1
      dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java
  42. 6
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java
  43. 4
      ui-ngx/src/app/modules/home/components/alarm-rules/alarm-rules-table-config.ts
  44. 2
      ui-ngx/src/app/modules/home/components/alarm-rules/cf-alarm-rule.component.html
  45. 3
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-complex-filter-predicate-dialog.component.ts
  46. 8
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-list.component.ts
  47. 2
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-list.component.html
  48. 2
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-list.component.ts
  49. 10
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-value.component.ts
  50. 14
      ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate.component.ts
  51. 1
      ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html
  52. 12
      ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/calculated-field-geofencing-zone-groups-panel.component.ts
  53. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/calculated-field-geofencing-zone-groups-table.component.ts
  54. 1
      ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/geofencing-configuration.component.html
  55. 3
      ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/geofencing-configuration.component.ts
  56. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts
  57. 2
      ui-ngx/src/app/modules/home/pages/asset/asset-tabs.component.html
  58. 12
      ui-ngx/src/app/modules/home/pages/device-profile/device-profile-tabs.component.html
  59. 2
      ui-ngx/src/app/modules/home/pages/device-profile/device-profile-tabs.component.ts
  60. 4
      ui-ngx/src/app/modules/home/pages/device/device-tabs.component.html
  61. 1
      ui-ngx/src/app/shared/models/calculated-field.models.ts
  62. 2
      ui-ngx/src/app/shared/models/tenant.model.ts
  63. 4
      ui-ngx/src/assets/locale/locale.constant-en_US.json

2
application/src/main/data/upgrade/basic/schema_update.sql

@ -45,7 +45,7 @@ SET profile_data = jsonb_set(
CASE
WHEN (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF'
THEN NULL
ELSE to_jsonb(3600)
ELSE to_jsonb(60)
END,
'minAggregationIntervalInSecForCF',
CASE

17
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -415,6 +415,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} catch (Exception e) {
throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build();
}
} else if (ctx.shouldFetchEntityRelations(state)) {
log.debug("[{}][{}] Going to update related entities for CF.", entityId, ctx.getCfId());
try {
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesState) {
List<EntityId> relatedEntities = cfService.fetchRelatedEntities(ctx, entityId);
List<EntityId> missingEntities = relatedEntitiesState.checkRelatedEntities(relatedEntities);
if (!missingEntities.isEmpty()) {
missingEntities.forEach(missingEntityId -> {
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, missingEntityId, ctx.getArguments());
relatedEntitiesState.updateEntityData(setEntityIdToSingleEntityArguments(missingEntityId, fetchedArgs));
});
justRestored = true;
}
}
} catch (Exception e) {
throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build();
}
}
if (state.isSizeOk()) {
Map<String, ArgumentEntry> updatedArgs = state.update(newArgValues, ctx);

18
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -27,9 +27,11 @@ import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.ProfileEntityIdInfo;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.AssetProfile;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
@ -74,6 +76,7 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ScheduledFuture;
@ -320,12 +323,16 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
EntityId fromId = entityRelation.getFrom();
String relationType = entityRelation.getType();
if (!(CalculatedField.isSupportedRefEntity(toId) || CalculatedField.isSupportedRefEntity(fromId))) {
callback.onSuccess();
return;
}
MultipleTbCallback callbackForToAndFrom = new MultipleTbCallback(2, callback);
processRelationByDirection(EntitySearchDirection.TO, relationType, toId, callbackForToAndFrom, relationAction.apply(fromId));
processRelationByDirection(EntitySearchDirection.FROM, relationType, fromId, callbackForToAndFrom, relationAction.apply(toId));
}
private void processRelationByDirection(EntitySearchDirection direction,
String relationType,
EntityId mainId,
@ -339,12 +346,11 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
List<CalculatedFieldCtx> matchingCfs = cfsByEntityIdAndProfile.stream()
.filter(cf -> {
if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config ) {
if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) {
RelationPathLevel relation = config.getRelation();
return direction.equals(relation.direction()) && relationType.equals(relation.relationType());
} else {
return false;
}
return false;
})
.toList();
@ -717,8 +723,8 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
private EntityId getProfileId(TenantId tenantId, EntityId entityId) {
return switch (entityId.getEntityType()) {
case ASSET -> assetProfileCache.get(tenantId, (AssetId) entityId).getId();
case DEVICE -> deviceProfileCache.get(tenantId, (DeviceId) entityId).getId();
case ASSET -> Optional.ofNullable(assetProfileCache.get(tenantId, (AssetId) entityId)).map(AssetProfile::getId).orElse(null);
case DEVICE -> Optional.ofNullable(deviceProfileCache.get(tenantId, (DeviceId) entityId)).map(DeviceProfile::getId).orElse(null);
default -> null;
};
}

25
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -24,6 +24,7 @@ import jakarta.annotation.PreDestroy;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration;
@ -62,6 +63,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION;
@ -184,11 +186,12 @@ public abstract class AbstractCalculatedFieldProcessingService {
}
protected Map<String, ListenableFuture<ArgumentEntry>> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) {
RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration();
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, aggConfig.getRelation());
if (!(ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config)) {
return Collections.emptyMap();
}
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, config.getRelation());
return aggConfig.getArguments().entrySet().stream()
return config.getArguments().entrySet().stream()
.collect(Collectors.toMap(
Map.Entry::getKey,
entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchRelatedEntitiesArgumentEntry(ctx.getTenantId(), relatedEntities, entry.getValue(), ts), MoreExecutors.directExecutor())
@ -205,8 +208,9 @@ public abstract class AbstractCalculatedFieldProcessingService {
));
}
private ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) {
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findByRelationPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation)));
protected ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) {
Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo());
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findFilteredRelationsByPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation)), filter);
return Futures.transform(relationsFut, relations -> {
if (relations == null) {
@ -217,7 +221,11 @@ public abstract class AbstractCalculatedFieldProcessingService {
case FROM -> relations.stream()
.map(EntityRelation::getTo)
.toList();
case TO -> relations.isEmpty() ? List.of() : List.of(relations.get(0).getFrom());
case TO -> relations.stream()
.map(EntityRelation::getFrom)
.findFirst()
.map(List::of)
.orElseGet(Collections::emptyList);
};
}, calculatedFieldCallbackExecutor);
}
@ -239,7 +247,8 @@ public abstract class AbstractCalculatedFieldProcessingService {
case CURRENT_OWNER -> Futures.immediateFuture(List.of(resolveOwnerArgument(tenantId, entityId)));
case RELATION_PATH_QUERY -> {
var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration;
yield Futures.transform(relationService.findByRelationPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId)),
Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo());
yield Futures.transform(relationService.findFilteredRelationsByPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId), filter),
configuration::resolveEntityIds, calculatedFieldCallbackExecutor);
}
};

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java

@ -37,6 +37,8 @@ public interface CalculatedFieldProcessingService {
Map<String, ArgumentEntry> fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId);
List<EntityId> fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId);
Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments);
ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval);

15
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric;
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -56,6 +57,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@ -99,6 +101,19 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
};
}
@Override
public List<EntityId> fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId) {
try {
if (ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) {
return resolveRelatedEntities(ctx.getTenantId(), entityId, config.getRelation()).get();
}
return Collections.emptyList();
} catch (ExecutionException | InterruptedException e) {
Throwable cause = e.getCause();
throw new RuntimeException("Failed to fetch related entities for entity [" + entityId + "]: " + cause.getMessage(), cause);
}
}
@Override
public Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments) {
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>();

3
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java

@ -162,7 +162,7 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS
}
private boolean checkEntityForCalculatedFields(TenantId tenantId, EntityId entityId, Predicate<CalculatedFieldCtx> filter, Predicate<CalculatedFieldCtx> linkedEntityFilter, Predicate<CalculatedFieldCtx> dynamicSourceFilter, Predicate<CalculatedFieldCtx> relatedEntityFilter) {
if (!CalculatedField.SUPPORTED_REFERENCED_ENTITIES.contains(entityId.getEntityType())) {
if (!CalculatedField.isSupportedRefEntity(entityId)) {
return false;
}
@ -211,7 +211,6 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS
return true;
}
}
return false;
}
}
}

14
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -61,6 +61,7 @@ import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.service.cf.CalculatedFieldProcessingService;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAggregationCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import org.thingsboard.server.service.telemetry.AlarmSubscriptionService;
@ -703,6 +704,19 @@ public class CalculatedFieldCtx implements Closeable {
};
}
public boolean shouldFetchEntityRelations(CalculatedFieldState state) {
if (!(state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState)) {
return false;
}
if (!isScheduledUpdateEnabled()) {
return false;
}
if (relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() == -1L) {
return true;
}
return relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis;
}
@Override
public void close() {
try {

87
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java

@ -38,9 +38,12 @@ import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.ScheduledFuture;
import static java.util.concurrent.TimeUnit.SECONDS;
@ -52,9 +55,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
private long lastArgsRefreshTs = -1;
@Setter
private long lastMetricsEvalTs = -1;
@Setter
private long lastRelatedEntitiesRefreshTs = -1;
private long deduplicationIntervalMs = -1;
private Map<String, AggMetric> metrics;
private ScheduledFuture<?> reevaluationFuture;
public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) {
super(entityId);
}
@ -67,8 +74,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec());
}
public void scheduleReevaluation() {
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
@Override
public void close() {
super.close();
if (reevaluationFuture != null) {
reevaluationFuture.cancel(true);
reevaluationFuture = null;
}
}
@Override
@ -76,9 +88,14 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
super.reset();
lastArgsRefreshTs = -1;
lastMetricsEvalTs = -1;
lastRelatedEntitiesRefreshTs = -1;
metrics = null;
}
public void updateLastRelatedEntitiesRefreshTs() {
lastRelatedEntitiesRefreshTs = System.currentTimeMillis();
}
@Override
public CalculatedFieldType getType() {
return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
@ -90,6 +107,56 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
return super.update(argumentValues, ctx);
}
public List<EntityId> checkRelatedEntities(List<EntityId> relatedEntities) {
Map<EntityId, Map<String, ArgumentEntry>> entityInputs = prepareInputs();
findOutdatedEntities(entityInputs, relatedEntities).forEach(this::cleanupEntityData);
updateLastRelatedEntitiesRefreshTs();
return findMissingEntities(entityInputs, relatedEntities);
}
private List<EntityId> findMissingEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) {
List<EntityId> missing = new ArrayList<>();
relatedEntities.forEach(entityId -> {
if (!entityInputs.containsKey(entityId)) {
missing.add(entityId);
log.warn("[{}] Missing related entity inputs for {}", ctx.getCfId(), entityId);
}
});
return missing;
}
private List<EntityId> findOutdatedEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) {
List<EntityId> outdated = new ArrayList<>();
entityInputs.keySet().forEach(entityId -> {
if (!relatedEntities.contains(entityId)) {
outdated.add(entityId);
log.warn("[{}] CF state keeps outdated related entity {}", ctx.getCfId(), entityId);
}
});
return outdated;
}
public Map<String, ArgumentEntry> updateEntityData(Map<String, ArgumentEntry> fetchedArgs) {
lastMetricsEvalTs = -1;
return update(fetchedArgs, ctx);
}
public void cleanupEntityData(EntityId relatedEntityId) {
arguments.values().forEach(argEntry -> {
RelatedEntitiesArgumentEntry aggEntry = (RelatedEntitiesArgumentEntry) argEntry;
aggEntry.getEntityInputs().remove(relatedEntityId);
});
lastMetricsEvalTs = -1;
lastArgsRefreshTs = System.currentTimeMillis();
}
public void scheduleReevaluation() {
ScheduledFuture<?> future = ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
if (future != null) {
reevaluationFuture = future;
}
}
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception {
boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty();
@ -97,7 +164,7 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
Output output = ctx.getOutput();
ObjectNode aggResult = aggregateMetrics(output);
lastMetricsEvalTs = System.currentTimeMillis();
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
scheduleReevaluation();
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.type(output.getType())
.scope(output.getScope())
@ -108,20 +175,6 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
}
}
public Map<String, ArgumentEntry> updateEntityData(Map<String, ArgumentEntry> fetchedArgs) {
lastMetricsEvalTs = -1;
return update(fetchedArgs, ctx);
}
public void cleanupEntityData(EntityId relatedEntityId) {
arguments.values().forEach(argEntry -> {
RelatedEntitiesArgumentEntry aggEntry = (RelatedEntitiesArgumentEntry) argEntry;
aggEntry.getEntityInputs().remove(relatedEntityId);
});
lastMetricsEvalTs = -1;
lastArgsRefreshTs = System.currentTimeMillis();
}
private boolean shouldRecalculate() {
boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationIntervalMs;
boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs;

13
application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java

@ -29,7 +29,6 @@ import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.ObjectType;
import org.thingsboard.server.common.data.TbResource;
import org.thingsboard.server.common.data.TbResourceInfo;
import org.thingsboard.server.common.data.Tenant;
@ -49,6 +48,7 @@ import org.thingsboard.server.common.data.job.Job;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.notification.NotificationRequest;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.data.security.DeviceCredentials;
@ -274,10 +274,13 @@ public class EntityStateSourcingListener {
@TransactionalEventListener(fallbackExecution = true)
public void handleEvent(RelationActionEvent relationEvent) {
if (relationEvent.getActionType() == ActionType.RELATION_ADD_OR_UPDATE) {
tbClusterService.onRelationUpdated(relationEvent.getTenantId(), relationEvent.getRelation(), TbQueueCallback.EMPTY);
} else if (relationEvent.getActionType() == ActionType.RELATION_DELETED) {
tbClusterService.onRelationDeleted(relationEvent.getTenantId(), relationEvent.getRelation(), TbQueueCallback.EMPTY);
EntityRelation relation = relationEvent.getRelation();
if (CalculatedField.isSupportedRefEntity(relation.getFrom()) && CalculatedField.isSupportedRefEntity(relation.getTo())) {
if (relationEvent.getActionType() == ActionType.RELATION_ADD_OR_UPDATE) {
tbClusterService.onRelationUpdated(relationEvent.getTenantId(), relation, TbQueueCallback.EMPTY);
} else if (relationEvent.getActionType() == ActionType.RELATION_DELETED) {
tbClusterService.onRelationDeleted(relationEvent.getTenantId(), relation, TbQueueCallback.EMPTY);
}
}
}

96
application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java

@ -51,6 +51,7 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationPathLevel;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
@ -87,6 +88,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
updateDefaultTenantProfileConfig(tenantProfileConfig -> {
tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1);
tenantProfileConfig.setMinAllowedScheduledUpdateIntervalInSecForCF(1);
});
Tenant tenant = new Tenant();
@ -177,7 +179,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
await().alias("add entity to profile with no related entities and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("add entity to profile with no related entities and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ObjectNode occupancy = getLatestTelemetry(asset2.getId(), "freeSpaces", "occupiedSpaces", "totalSpaces");
@ -190,7 +192,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), device4.getId(), "Contains");
await().alias("create relations and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create relations and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -202,7 +204,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device3.getId(), "{\"occupied\":false}");
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("update telemetry and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -224,7 +226,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createOccupancyCF(assetProfile.getId());
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -246,7 +248,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device3.getId(), "{\"occupied\":true}");
await().alias("change profile and no aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("change profile and no aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -268,7 +270,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createOccupancyCF(asset2.getId());
await().alias("create CF and perform aggregation with default values").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create CF and perform aggregation with default values").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -299,6 +301,45 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
});
}
@Test
public void testCreateCfAndRelationToRuleChain_checkAggregation() throws Exception {
Asset asset2 = createAsset("Asset 2", assetProfile.getId());
Device device3 = createDevice("Device 3", "1234567890333");
postTelemetry(device3.getId(), "{\"occupied\":true}");
RuleChain ruleChain = new RuleChain();
ruleChain.setName("RuleChain");
ruleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class);
postTelemetry(ruleChain.getId(), "{\"occupied\":true}");
createEntityRelation(asset2.getId(), device3.getId(), "Contains");
createEntityRelation(asset2.getId(), ruleChain.getId(), "Contains");
createOccupancyCF(asset2.getId());
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "0",
"occupiedSpaces", "1",
"totalSpaces", "1"
));
});
postTelemetry(ruleChain.getId(), "{\"occupied\":true}");
await().alias("update telemetry on rule chain and no aggregation performed").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
"freeSpaces", "0",
"occupiedSpaces", "1",
"totalSpaces", "1"
));
});
}
@Test
public void testDeleteCf_checkNoAggregation() throws Exception {
CalculatedField cf = createOccupancyCF(asset.getId());
@ -309,7 +350,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device1.getId(), "{\"occupied\":false}");
await().alias("delete cf and update telemetry and no aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("delete cf and update telemetry and no aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -364,7 +405,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createOccupancyCF(asset2.getId());
await().alias("create CF and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create CF and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -402,7 +443,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createOccupancyCFWithAttr(asset2.getId());
await().alias("create CF and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create CF and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset2.getId(), Map.of(
@ -437,7 +478,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
createEntityRelation(asset.getId(), device3.getId(), "Contains");
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create relation and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -455,7 +496,25 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
deleteEntityRelation(new EntityRelation(asset.getId(), device1.getId(), "Contains", RelationTypeGroup.COMMON));
await().alias("create relation and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create relation and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
"freeSpaces", "1",
"occupiedSpaces", "0",
"totalSpaces", "1"
));
});
}
@Test
public void testDeleteEntityByRelation_checkAggregation() throws Exception {
createOccupancyCF(asset.getId());
checkInitialCalculation();
doDelete("/api/device/" + device1.getId()).andExpect(status().isOk());
await().alias("create relation and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -479,7 +538,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
configuration.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, "Has"));
saveCalculatedField(cf);
await().alias("update relation path and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("update relation path and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -505,7 +564,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
configuration.setArguments(Map.of("oc", argument));
saveCalculatedField(cf);
await().alias("update arguments and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("update arguments and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of(
@ -560,7 +619,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device2.getId(), "{\"temperature\":19.6}");
CalculatedField cf = createAvgTemperatureCF(asset.getId());
await().alias("create avg temp cf and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create avg temp cf and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of("avgTemperature", "24"));
@ -573,7 +632,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
configuration.setOutput(output);
saveCalculatedField(cf);
await().alias("update output and perform aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("update output and perform aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
ArrayNode avgTemperature = getServerAttributes(asset.getId(), "avgTemperature");
@ -589,7 +648,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device2.getId(), "{\"temperature\":19.6}");
CalculatedField cf = createAvgTemperatureCF(asset.getId());
await().alias("create avg temp cf and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create avg temp cf and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of("avgTemperature", "24"));
@ -607,7 +666,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
postTelemetry(device2.getId(), "{\"temperature\":32.1}");
await().alias("update telemetry and perform aggregation").atMost(2 * deduplicationInterval, TimeUnit.SECONDS)
await().alias("update telemetry and perform aggregation").atMost(2 * deduplicationInterval + 10, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(() -> {
verifyTelemetry(asset.getId(), Map.of("avgTemperature", "28"));
@ -615,7 +674,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
}
private void checkInitialCalculation() {
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval, TimeUnit.SECONDS)
await().alias("create CF and perform initial aggregation").atMost(deduplicationInterval * 2, TimeUnit.SECONDS)
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS)
.untilAsserted(this::checkInitialCalculationValues);
}
@ -743,6 +802,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
configuration.setRelation(relation);
configuration.setArguments(inputs);
configuration.setDeduplicationIntervalInSec(deduplicationInterval);
configuration.setScheduledUpdateInterval(10);
configuration.setMetrics(metrics);
configuration.setOutput(output);

27
application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java

@ -128,11 +128,17 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest {
protected void testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(HasName entity, EntityId entityId, EntityId originatorId,
TenantId tenantId, CustomerId customerId, UserId userId, String userName,
ActionType actionType, ActionType actionTypeEdge, Object... additionalInfo) {
testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(tenantId, entity, entityId, originatorId, tenantId, customerId, userId, userName, actionType, actionTypeEdge, additionalInfo);
}
protected void testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(TenantId entityTenantId, HasName entity, EntityId entityId, EntityId originatorId,
TenantId authTenantId, CustomerId customerId, UserId userId, String userName,
ActionType actionType, ActionType actionTypeEdge, Object... additionalInfo) {
int cntTime = 1;
testNotificationMsgToEdgeServiceTime(entityId, tenantId, actionTypeEdge, cntTime);
testLogEntityActionEntityEqClass(entity, originatorId, tenantId, customerId, userId, userName, actionType, cntTime, additionalInfo);
testNotificationMsgToEdgeServiceTime(entityId, entityTenantId, actionTypeEdge, cntTime);
testLogEntityActionEntityEqClass(entity, originatorId, authTenantId, customerId, userId, userName, actionType, cntTime, additionalInfo);
ArgumentMatcher<EntityId> matcherOriginatorId = argument -> argument.equals(originatorId);
testPushMsgToRuleEngineTime(matcherOriginatorId, tenantId, entity, cntTime);
testPushMsgToRuleEngineTime(matcherOriginatorId, authTenantId, entity, cntTime);
Mockito.reset(tbClusterService, auditLogService);
}
@ -159,17 +165,26 @@ public abstract class AbstractNotifyEntityTest extends AbstractWebTest {
TenantId tenantId, CustomerId customerId, UserId userId, String userName,
ActionType actionType,
int cntTime, int cntTimeEdge, int cntTimeRuleEngine, Object... additionalInfo) {
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(tenantId, entity, originator, tenantId, customerId, userId, userName, actionType,
cntTime, cntTimeEdge, cntTimeRuleEngine, additionalInfo);
}
protected void testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(TenantId entityTenantId, HasName entity, HasName originator,
TenantId authTenantId, CustomerId customerId, UserId userId, String userName,
ActionType actionType,
int cntTime, int cntTimeEdge, int cntTimeRuleEngine, Object... additionalInfo) {
EntityId originatorId = createEntityId_NULL_UUID(originator);
testSendNotificationMsgToEdgeServiceTimeEntityEqAny(tenantId, actionType, cntTimeEdge);
testSendNotificationMsgToEdgeServiceTimeEntityEqAny(entityTenantId, actionType, cntTimeEdge);
ArgumentMatcher<HasName> matcherEntityClassEquals = argument -> argument.getClass().equals(entity.getClass());
ArgumentMatcher<EntityId> matcherOriginatorId = argument -> argument.getClass().equals(originatorId.getClass());
ArgumentMatcher<CustomerId> matcherCustomerId = customerId == null ?
argument -> argument.getClass().equals(CustomerId.class) : argument -> argument.equals(customerId);
ArgumentMatcher<UserId> matcherUserId = userId == null ?
argument -> argument.getClass().equals(UserId.class) : argument -> argument.equals(userId);
testLogEntityActionAdditionalInfo(matcherEntityClassEquals, matcherOriginatorId, tenantId, matcherCustomerId, matcherUserId, userName, actionType, cntTime,
testLogEntityActionAdditionalInfo(matcherEntityClassEquals, matcherOriginatorId, authTenantId, matcherCustomerId, matcherUserId, userName, actionType, cntTime,
extractMatcherAdditionalInfoClass(additionalInfo));
testPushMsgToRuleEngineTime(matcherOriginatorId, tenantId, entity, cntTimeRuleEngine);
testPushMsgToRuleEngineTime(matcherOriginatorId, authTenantId, entity, cntTimeRuleEngine);
}
protected void testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAnyAdditionalInfoAny(HasName entity, HasName originator,

2
application/src/test/java/org/thingsboard/server/controller/TbResourceControllerTest.java

@ -943,7 +943,7 @@ public class TbResourceControllerTest extends AbstractControllerTest {
private List<TbResourceInfo> loadLwm2mResources() throws Exception {
var models = List.of("1", "2", "3", "5", "6", "9", "19", "3303");
var models = List.of("1", "2", "3-1_2", "5", "6", "9", "19", "3303");
List<TbResourceInfo> resources = new ArrayList<>(models.size());

8
application/src/test/java/org/thingsboard/server/controller/UserControllerTest.java

@ -116,7 +116,7 @@ public class UserControllerTest extends AbstractControllerTest {
foundUser.setAdditionalInfo(savedUser.getAdditionalInfo());
Assert.assertEquals(foundUser, savedUser);
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(foundUser, foundUser,
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(user.getTenantId(), foundUser, foundUser,
SYSTEM_TENANT, customerNUULId, null, SYS_ADMIN_EMAIL,
ActionType.ADDED, 1, 1, 1);
Mockito.reset(tbClusterService, auditLogService);
@ -155,7 +155,7 @@ public class UserControllerTest extends AbstractControllerTest {
doDelete("/api/user/" + savedUser.getId().getId().toString())
.andExpect(status().isOk());
testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(foundUser, foundUser.getId(), foundUser.getId(),
testNotifyEntityAllOneTimeLogEntityActionEntityEqClass(user.getTenantId(), foundUser, foundUser.getId(), foundUser.getId(),
SYSTEM_TENANT, customerNUULId, null, SYS_ADMIN_EMAIL,
ActionType.DELETED, ActionType.DELETED, SYSTEM_TENANT.getId().toString());
}
@ -414,7 +414,7 @@ public class UserControllerTest extends AbstractControllerTest {
User testManyUser = new User();
testManyUser.setTenantId(tenantId);
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(testManyUser, testManyUser,
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(tenantId, testManyUser, testManyUser,
SYSTEM_TENANT, customerNUULId, null, SYS_ADMIN_EMAIL,
ActionType.ADDED, cntEntity, cntEntity, cntEntity);
@ -526,7 +526,7 @@ public class UserControllerTest extends AbstractControllerTest {
}
User testManyUser = new User();
testManyUser.setTenantId(tenantId);
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(testManyUser, testManyUser,
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(tenantId, testManyUser, testManyUser,
SYSTEM_TENANT, customerNUULId, null, SYS_ADMIN_EMAIL,
ActionType.DELETED, cntEntity, NUMBER_OF_USERS, cntEntity, "");

38
application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java

@ -21,6 +21,7 @@ import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.io.IOUtils;
import org.awaitility.core.ConditionTimeoutException;
import org.eclipse.leshan.client.LeshanClient;
import org.eclipse.leshan.client.object.Security;
import org.eclipse.leshan.client.servers.LwM2mServer;
@ -95,6 +96,7 @@ import java.util.Map;
import java.util.Set;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.awaitility.Awaitility.await;
import static org.eclipse.leshan.client.object.Security.noSec;
@ -119,6 +121,7 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClient
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_UPDATE_SUCCESS;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.lwm2mClientResources;
import static org.thingsboard.server.transport.lwm2m.ota.AbstractOtaLwM2MIntegrationTest.CLIENT_LWM2M_SETTINGS_19;
@Slf4j
@ -304,7 +307,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
protected final Set<Lwm2mTestHelper.LwM2MClientState> expectedStatusesRegistrationBsSuccess = new HashSet<>(Arrays.asList(ON_BOOTSTRAP_STARTED, ON_BOOTSTRAP_SUCCESS, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS));
protected ScheduledExecutorService executor;
protected LwM2MTestClient lwM2MTestClient;
private String[] resources;
private String[] resources = lwm2mClientResources;
protected String deviceId;
protected boolean supportFormatOnly_SenMLJSON_SenMLCBOR = false;
@ -546,7 +549,9 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
}
public void setResources(String[] resources) {
this.resources = resources;
if (this.resources == null || !Arrays.equals(this.resources, resources)) {
this.resources = resources;
}
}
public void createNewClient(Security security, Security securityBs, boolean isRpc,
@ -741,11 +746,19 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
return credentials;
}
protected void awaitObserveReadAll(int cntObserve, String deviceIdStr) throws Exception {
await("ObserveReadAll: countObserve " + cntObserve)
.atMost(40, TimeUnit.SECONDS)
.until(() -> cntObserve == getCntObserveAll(deviceIdStr));
protected void awaitObserveReadAll(int cntObserve, String deviceIdStr) throws Exception {
try {
await("ObserveReadAll: countObserve " + cntObserve)
.atMost(40, TimeUnit.SECONDS)
.until(() -> cntObserve == getCntObserveAll(deviceIdStr));
} catch (ConditionTimeoutException e) {
int current = getCntObserveAll(deviceIdStr);
log.error("Condition or device {} with alias 'ObserveReadAll: countObserve {}, but received {}", deviceIdStr, cntObserve, current);
throw e;
}
}
protected void awaitDeleteDevice(String deviceIdStr) throws Exception {
await("Delete device with id: " + deviceIdStr)
.atMost(40, TimeUnit.SECONDS)
@ -756,6 +769,19 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
});
}
protected void updateRegAtLeastOnceAfterAction() {
long initialInvocationCount = countUpdateReg();
AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount);
log.trace("updateRegAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount);
await("Update Registration at-least-once after action")
.atMost(50, TimeUnit.SECONDS)
.until(() -> {
newInvocationCount.set(countUpdateReg());
return newInvocationCount.get() > initialInvocationCount;
});
log.trace("updateRegAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get());
}
protected Integer getCntObserveAll(String deviceIdStr) throws Exception {
String actualResult = sendObserveOK("ObserveReadAll", null, deviceIdStr);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);

2
application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java

@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m;
public class Lwm2mTestHelper {
public static final String[] lwm2mClientResources = new String[]{"3.xml", "5.xml", "6.xml", "9.xml", "19.xml", "3303.xml"};
public static final String[] lwm2mClientResources = new String[]{"3-1_2.xml", "5.xml", "6.xml", "9.xml", "19.xml", "3303.xml"};
// Models
public static final int BINARY_APP_DATA_CONTAINER = 19;

36
application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java

@ -16,6 +16,7 @@
package org.thingsboard.server.transport.lwm2m.client;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.client.LeshanClient;
import org.eclipse.leshan.client.resource.BaseInstanceEnabler;
import org.eclipse.leshan.client.servers.LwM2mServer;
import org.eclipse.leshan.core.model.ObjectModel;
@ -32,6 +33,9 @@ import java.util.List;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicInteger;
import static org.thingsboard.server.dao.service.OtaPackageServiceTest.TARGET_FW_VERSION;
import static org.thingsboard.server.dao.service.OtaPackageServiceTest.TITLE;
@Slf4j
public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
@ -44,6 +48,12 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
private final AtomicInteger updateResult = new AtomicInteger(0);
private LeshanClient leshanClient;
private String pkgNameDef = "firmware";
private String pkgName;
private String pkgVersionDef = "1.0.0";
private String pkgVersion;
@Override
public ReadResponse read(LwM2mServer identity, int resourceId) {
if (!identity.isSystem())
@ -74,7 +84,7 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
switch (resourceId) {
case 2:
startUpdating();
startUpdating(identity);
return ExecuteResponse.success();
default:
return super.execute(identity, resourceId, arguments);
@ -106,11 +116,13 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
}
private String getPkgName() {
return "firmware";
this.pkgName = this.pkgName == null ? this.pkgNameDef : this.pkgName;
return this.pkgName;
}
private String getPkgVersion() {
return "1.0.0";
this.pkgVersion = this.pkgVersion == null ? this.pkgVersionDef : this.pkgVersion;
return this.pkgVersion;
}
private int getFirmwareUpdateDeliveryMethod() {
@ -140,7 +152,7 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
}, 100, TimeUnit.MILLISECONDS);
}
private void startUpdating() {
private void startUpdating(LwM2mServer identity) {
scheduler.schedule(() -> {
try {
state.set(3);
@ -148,9 +160,25 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
Thread.sleep(100);
updateResult.set(1);
fireResourceChange(5);
this.pkgName = TITLE;
fireResourceChange(6);
this.pkgVersion = TARGET_FW_VERSION;
fireResourceChange(7);
if (this.leshanClient != null) {
log.info("Stop/reboot LwM2M client {}", this.leshanClient.getEndpoint(identity));
this.leshanClient.stop(false);
log.info("Start after update fw LwM2M client {}", this.leshanClient.getEndpoint(identity));
this.leshanClient.start();
this.pkgName = this.pkgNameDef;
this.pkgVersion = this.pkgVersionDef;
}
} catch (Exception e) {
}
}, 100, TimeUnit.MILLISECONDS);
}
protected void setLeshanClient(LeshanClient leshanClient) {
this.leshanClient = leshanClient;
}
}

1
application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java

@ -467,6 +467,7 @@ public class LwM2MTestClient {
this.awaitClientAfterStartConnectLw();
}
lwM2mTemperatureSensor12.setLeshanClient(leshanClient);
fwLwM2MDevice.setLeshanClient(leshanClient);
}
}

5
application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java

@ -54,7 +54,6 @@ import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaU
@DaoSqlTest
public abstract class AbstractOtaLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest {
private final String[] RESOURCES_OTA = new String[]{"3.xml", "5.xml", "9.xml", "19.xml"};
protected static final String CLIENT_ENDPOINT_WITHOUT_FW_INFO = "WithoutFirmwareInfoDevice";
protected static final String CLIENT_ENDPOINT_OTA5 = "Ota5_Device";
protected static final String CLIENT_ENDPOINT_OTA9 = "Ota9_Device";
@ -186,10 +185,6 @@ public abstract class AbstractOtaLwM2MIntegrationTest extends AbstractLwM2MInteg
" \"attributeLwm2m\": {}\n" +
" }";
public AbstractOtaLwM2MIntegrationTest() {
setResources(this.RESOURCES_OTA);
}
protected OtaPackageInfo createFirmware(String version, DeviceProfileId deviceProfileId) throws Exception {
String CHECKSUM = "4bf5122f344554c53bde2ebb8cd2b7e3d1600ad631c385a5d7cce23c7785459a";

26
application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/Ota5LwM2MIntegrationTest.java

@ -45,6 +45,7 @@ import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.INIT
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.QUEUED;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.UPDATED;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.UPDATING;
import static org.thingsboard.server.dao.service.OtaPackageServiceTest.TARGET_FW_VERSION;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.BINARY_APP_DATA_CONTAINER;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_0;
@ -91,13 +92,14 @@ public class Ota5LwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest {
@Test
public void testFirmwareUpdateByObject5_Ok() throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(OBSERVE_ATTRIBUTES_WITH_PARAMS_OTA5, getBootstrapServerCredentialsNoSec(NONE));
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + this.CLIENT_ENDPOINT_OTA5, transportConfiguration);
LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(this.CLIENT_ENDPOINT_OTA5));
final Device device = createLwm2mDevice(deviceCredentials, this.CLIENT_ENDPOINT_OTA5, deviceProfile.getId());
createNewClient(SECURITY_NO_SEC, null, false, this.CLIENT_ENDPOINT_OTA5, device.getId().getId().toString());
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + this.CLIENT_ENDPOINT_OTA5 + "Ok", transportConfiguration);
String endpoint = this.CLIENT_ENDPOINT_OTA5 + "Ok";
LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(endpoint));
final Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId());
createNewClient(SECURITY_NO_SEC, null, false, endpoint, device.getId().getId().toString());
awaitObserveReadAll(5, device.getId().getId().toString());
device.setFirmwareId(createFirmware("fw.v.1.5.0-update", deviceProfile.getId()).getId());
device.setFirmwareId(createFirmware(TARGET_FW_VERSION, deviceProfile.getId()).getId());
final Device savedDevice = doPost("/api/device", device, Device.class);
assertThat(savedDevice).as("saved device").isNotNull();
@ -110,7 +112,6 @@ public class Ota5LwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest {
log.warn("Object5: Got the ts: {}", ts);
}
/**
* ObjectId = 19/65533/0
* {
@ -133,13 +134,14 @@ public class Ota5LwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest {
@Test
public void testFirmwareUpdateByObject5WithObject19_Ok() throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration19(OBSERVE_ATTRIBUTES_WITH_PARAMS_OTA5_19, getBootstrapServerCredentialsNoSec(NONE));
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + this.CLIENT_ENDPOINT_OTA5, transportConfiguration);
LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(this.CLIENT_ENDPOINT_OTA5));
final Device device = createLwm2mDevice(deviceCredentials, this.CLIENT_ENDPOINT_OTA5, deviceProfile.getId());
createNewClient(SECURITY_NO_SEC, null, false, this.CLIENT_ENDPOINT_OTA5, device.getId().getId().toString());
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + this.CLIENT_ENDPOINT_OTA5 + "19_Ok", transportConfiguration);
String endpoint = this.CLIENT_ENDPOINT_OTA5 + "19_Ok";
LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(endpoint));
final Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId());
createNewClient(SECURITY_NO_SEC, null, false, endpoint, device.getId().getId().toString());
awaitObserveReadAll(6, device.getId().getId().toString());
OtaPackageInfo otaPackageInfo = createFirmware("fw.v.1.5.0-update", deviceProfile.getId());
OtaPackageInfo otaPackageInfo = createFirmware(TARGET_FW_VERSION, deviceProfile.getId());
device.setFirmwareId(otaPackageInfo.getId());
final Device savedDevice = doPost("/api/device", device, Device.class);
@ -154,6 +156,6 @@ public class Ota5LwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest {
String ver_Id_19 = lwM2MTestClient.getLeshanClient().getObjectTree().getModel().getObjectModel(BINARY_APP_DATA_CONTAINER).version;
String resourceIdVer = "/" + BINARY_APP_DATA_CONTAINER + "_" + ver_Id_19 + "/" + FW_INSTANCE_ID + "/" + RESOURCE_ID_0;
resultReadOtaParams_19(resourceIdVer, otaPackageInfo);
log.warn("Object5: Got the ts: {}", ts);
log.warn("Object5 with Object19: Got the ts: {}", ts);
}
}

7
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test.java

@ -20,9 +20,8 @@ import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
public abstract class AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test extends AbstractRpcLwM2MIntegrationTest{
public AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test() {
String[] RESOURCES_RPC_VER_1_1 = new String[]{"3-1_0.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_RPC_VER_1_1);
public AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test() throws Exception {
String[] RESOURCES_RPC_VER_1_0 = new String[]{"3-1_0.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_RPC_VER_1_0);
}
}

2
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test.java

@ -20,7 +20,7 @@ import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
public abstract class AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test extends AbstractRpcLwM2MIntegrationTest{
public AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test() {
public AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test() throws Exception {
String[] RESOURCES_RPC_VER_1_1 = new String[]{"3-1_1.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_RPC_VER_1_1);
}

6
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test.java

@ -20,9 +20,9 @@ import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
public abstract class AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test extends AbstractRpcLwM2MIntegrationTest{
public AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test() {
String[] RESOURCES_RPC_VER_1_1 = new String[]{"3.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_RPC_VER_1_1);
public AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test() throws Exception {
String[] RESOURCES_RPC_VER_1_2 = new String[]{"3-1_2.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_RPC_VER_1_2);
}
}

54
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java

@ -15,7 +15,9 @@
*/
package org.thingsboard.server.transport.lwm2m.rpc;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.ResponseCode;
import org.eclipse.leshan.core.link.LinkParser;
import org.eclipse.leshan.core.link.lwm2m.DefaultLwM2mLinkParser;
import org.junit.Before;
@ -40,6 +42,9 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Predicate;
import static org.awaitility.Awaitility.await;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.eclipse.leshan.core.LwM2mId.ACCESS_CONTROL;
import static org.eclipse.leshan.core.LwM2mId.DEVICE;
import static org.eclipse.leshan.core.LwM2mId.FIRMWARE;
@ -101,10 +106,6 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg
@SpyBean
protected LwM2mTransportServerHelper lwM2mTransportServerHelperTest;
public AbstractRpcLwM2MIntegrationTest() {
setResources(lwm2mClientResources);
}
@Before
public void startInitRPC() throws Exception {
if (this.getClass().getSimpleName().equals("RpcLwm2mIntegrationWriteCborTest")) {
@ -324,19 +325,6 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg
.count();
}
protected void updateRegAtLeastOnceAfterAction() {
long initialInvocationCount = countUpdateReg();
AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount);
log.trace("updateRegAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount);
await("Update Registration at-least-once after action")
.atMost(50, TimeUnit.SECONDS)
.until(() -> {
newInvocationCount.set(countUpdateReg());
return newInvocationCount.get() > initialInvocationCount;
});
log.trace("updateRegAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get());
}
protected long countSendParametersOnThingsboardTelemetryResource(String rezName) {
return Mockito.mockingDetails(lwM2mTransportServerHelperTest)
.getInvocations().stream()
@ -350,4 +338,36 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg
)
.count();
}
protected String sendDiscover(String path) throws Exception {
String setRpcRequest = "{\"method\": \"Discover\", \"params\": {\"id\": \"" + path + "\"}}";
return doPostAsync("/api/plugins/rpc/twoway/" + lwM2MTestClient.getDeviceIdStr(), setRpcRequest, String.class, status().isOk());
}
protected String sendRpcObserveReadAllWithResult() throws Exception {
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveReadAll", null);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
protected String sendRpcObserveReadAllWithResult(String params) throws Exception {
sendRpcObserveOk("Observe", params);
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveReadAll", null);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
protected void testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4(String expectedIdVer) throws Exception {
String expectedIdObserve = "SingleObservation:/3/0/9";
sendObserveCancelAllWithAwait(lwM2MTestClient.getDeviceIdStr());
updateRegAtLeastOnceAfterAction();
lwM2MTestClient.getLeshanClient().stop(false);
lwM2MTestClient.getLeshanClient().start();
updateRegAtLeastOnceAfterAction();
awaitObserveReadAll(4,lwM2MTestClient.getDeviceIdStr());
String actualIdVer = sendDiscover(objectIdVer_3);
assertTrue(actualIdVer.contains(expectedIdVer));
String actualAllObserve = sendRpcObserveReadAllWithResult();
assertTrue(actualAllObserve.contains(expectedIdObserve));
}
}

5
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationDiscoverTest.java

@ -192,11 +192,6 @@ public class RpcLwm2mIntegrationDiscoverTest extends AbstractRpcLwM2MIntegration
assertTrue(rpcActualResult.get("error").asText().contains(expected));
}
private String sendDiscover(String path) throws Exception {
String setRpcRequest = "{\"method\": \"Discover\", \"params\": {\"id\": \"" + path + "\"}}";
return doPostAsync("/api/plugins/rpc/twoway/" + lwM2MTestClient.getDeviceIdStr(), setRpcRequest, String.class, status().isOk());
}
private String convertObjectIdToVerId(String path, String ver) {
ver = ver != null ? ver : TbLwM2mVersion.VERSION_1_0.getVersion().toString();
try {

5
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationDiscoverWriteAttributesTest.java

@ -166,9 +166,4 @@ public class RpcLwm2mIntegrationDiscoverWriteAttributesTest extends AbstractRpcL
String setRpcRequest = "{\"method\": \"WriteAttributes\", \"params\": {\"id\": \"" + path + "\", \"attributes\": " + value + " }}";
return doPostAsync("/api/plugins/rpc/twoway/" + lwM2MTestClient.getDeviceIdStr(), setRpcRequest, String.class, status().isOk());
}
private String sendDiscover(String path) throws Exception {
String setRpcRequest = "{\"method\": \"Discover\", \"params\": {\"id\": \"" + path + "\"}}";
return doPostAsync("/api/plugins/rpc/twoway/" + lwM2MTestClient.getDeviceIdStr(), setRpcRequest, String.class, status().isOk());
}
}

7
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java

@ -335,12 +335,5 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT
sendRpcObserveOk("Observe", expectedId_1);
sendRpcObserveOk("Observe", expectedId_2);
}
private String sendRpcObserveReadAllWithResult(String params) throws Exception {
sendRpcObserveOk("Observe", params);
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveReadAll", null);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
}

23
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserve_Ver_1_0_Test.java → application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer10Test.java

@ -19,13 +19,14 @@ import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test;
import static org.junit.Assert.assertTrue;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9;
@Slf4j
public class RpcLwm2mIntegrationObserve_Ver_1_0_Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test {
public class RpcLwm2mIntegrationObserveVer10Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test {
public RpcLwm2mIntegrationObserveVer10Test() throws Exception {
}
@Before
public void setupObserveTest() throws Exception {
@ -44,5 +45,21 @@ public class RpcLwm2mIntegrationObserve_Ver_1_0_Test extends AbstractRpcLwM2MInt
updateRegAtLeastOnceAfterAction();
long lastSendTelemetryAtCount = countSendParametersOnThingsboardTelemetryResource(RESOURCE_ID_NAME_3_9);
assertTrue(lastSendTelemetryAtCount > initSendTelemetryAtCount);
awaitObserveReadAll(1,lwM2MTestClient.getDeviceIdStr());
}
/**
* "3_1.0/0/9"
* Observe count 4
* CancelAll Observe
* Reboot
* Observe count 4 contains
* "/3_1.0" - Discover Object - find ver
* @throws Exception
*/
@Test
public void testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4_ObjectVer_1_0() throws Exception {
String expectedIdVer = "</3>;ver=1.0";
testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4(expectedIdVer);
}
}

20
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserve_Ver_1_1_Test.java → application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer11Test.java

@ -23,7 +23,10 @@ import static org.junit.Assert.assertTrue;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9;
@Slf4j
public class RpcLwm2mIntegrationObserve_Ver_1_1_Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test {
public class RpcLwm2mIntegrationObserveVer11Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_1_Test {
public RpcLwm2mIntegrationObserveVer11Test() throws Exception {
}
@Before
public void setupObserveTest() throws Exception {
@ -43,4 +46,19 @@ public class RpcLwm2mIntegrationObserve_Ver_1_1_Test extends AbstractRpcLwM2MInt
long lastSendTelemetryAtCount = countSendParametersOnThingsboardTelemetryResource(RESOURCE_ID_NAME_3_9);
assertTrue(lastSendTelemetryAtCount > initSendTelemetryAtCount);
}
/**
* "3_1.1/0/9"
* Observe count 4
* CancelAll Observe
* Reboot
* Observe count 4 contains
* "/3" - Discover Object - find ver (lwm2mVersion == 1.1)
* @throws Exception
*/
@Test
public void testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4_ObjectVer_1_1() throws Exception {
String expectedIdVer = "</3>";
testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4(expectedIdVer);
}
}

21
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserve_Ver_1_2_Test.java → application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveVer12Test.java

@ -18,14 +18,16 @@ package org.thingsboard.server.transport.lwm2m.rpc.sql;
import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserve_Ver_1_0_Test;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test;
import static org.junit.Assert.assertTrue;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9;
@Slf4j
public class RpcLwm2mIntegrationObserve_Ver_1_2_Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test {
public class RpcLwm2mIntegrationObserveVer12Test extends AbstractRpcLwM2MIntegrationObserve_Ver_1_2_Test {
public RpcLwm2mIntegrationObserveVer12Test() throws Exception {
}
@Before
public void setupObserveTest() throws Exception {
@ -45,4 +47,19 @@ public class RpcLwm2mIntegrationObserve_Ver_1_2_Test extends AbstractRpcLwM2MInt
long lastSendTelemetryAtCount = countSendParametersOnThingsboardTelemetryResource(RESOURCE_ID_NAME_3_9);
assertTrue(lastSendTelemetryAtCount > initSendTelemetryAtCount);
}
/**
* "3_1.2/0/9"
* Observe count 4
* CancelAll Observe
* Reboot
* Observe count 4 contains
* "/3_1.2" - Discover Object - find ver
* @throws Exception
*/
@Test
public void testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4_ObjectVer_1_2() throws Exception {
String expectedIdVer = "</3>;ver=1.2";
testObserveOneResourceValue_Count_4_CancelAll_Reboot_After_Observe_Count_4(expectedIdVer);
}
}

9
application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java

@ -22,6 +22,7 @@ import org.eclipse.leshan.client.object.Security;
import org.eclipse.leshan.core.ResponseCode;
import org.eclipse.leshan.core.util.Hex;
import org.junit.Assert;
import org.junit.Before;
import org.springframework.test.web.servlet.MvcResult;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
@ -119,7 +120,6 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M
protected final PrivateKey clientPrivateKeyFromCertTrust; // client private key used for X509 and RPK
protected final X509Certificate clientX509CertTrustNo; // client certificate signed by intermediate, rootCA with a good CN ("host name")
protected final PrivateKey clientPrivateKeyFromCertTrustNo; // client private key used for X509 and RPK
private final String[] RESOURCES_SECURITY = new String[]{"1.xml", "2.xml", "3.xml", "5.xml", "9.xml", "19.xml"};
private final LwM2MBootstrapClientCredentials defaultBootstrapCredentials;
@ -134,7 +134,6 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M
public AbstractSecurityLwM2MIntegrationTest() {
// create client credentials
setResources(this.RESOURCES_SECURITY);
try {
// Get certificates from key store
char[] clientKeyStorePwd = CLIENT_STORE_PWD.toCharArray();
@ -178,6 +177,12 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M
defaultBootstrapCredentials.setLwm2mServer(serverCredentials);
}
@Before
public void init() throws Exception {
String[] RESOURCES_SECURITY = new String[]{"3-1_2.xml", "5.xml", "6.xml", "9.xml", "19.xml"};
setResources(RESOURCES_SECURITY);
}
public void basicTestConnectionStartBS(String clientEndpoint,
String awaitAlias,
LwM2MProfileBootstrapConfigType type,

14
application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/PskLwm2mIntegrationTest.java

@ -68,10 +68,12 @@ public class PskLwm2mIntegrationTest extends AbstractSecurityLwM2MIntegrationTes
ON_REGISTRATION_SUCCESS,
true);
}
@Test
public void testWithPskConnectLwm2mOneObserveSuccessUpdateProfileManyObserveUpdateRegistrationSuccess() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_PSK;
String identity = CLIENT_PSK_IDENTITY;
String suf = "UpdateReg";
String clientEndpoint = CLIENT_ENDPOINT_PSK + "_" + suf;
String identity = CLIENT_PSK_IDENTITY + "_" + suf;
String keyPsk = CLIENT_PSK_KEY;
PSKClientCredential clientCredentials = new PSKClientCredential();
clientCredentials.setEndpoint(clientEndpoint);
@ -103,10 +105,12 @@ public class PskLwm2mIntegrationTest extends AbstractSecurityLwM2MIntegrationTes
awaitObserveReadAll(2, lwm2mDevice.getId().getId().toString());
awaitUpdateReg(3);
}
@Test
public void testWithPskConnectLwm2mSuccessObserveSuccessUnRegClientUpdateProfileObserveConnectLwm2mSuccessOWithNewObserve() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_PSK;
String identity = CLIENT_PSK_IDENTITY;
String suf = "UnReg";
String clientEndpoint = CLIENT_ENDPOINT_PSK + "_" + suf;
String identity = CLIENT_PSK_IDENTITY + "_" + suf;
String keyPsk = CLIENT_PSK_KEY;
PSKClientCredential clientCredentials = new PSKClientCredential();
clientCredentials.setEndpoint(clientEndpoint);
@ -139,7 +143,7 @@ public class PskLwm2mIntegrationTest extends AbstractSecurityLwM2MIntegrationTes
Assert.assertNotNull(lwm2mDeviceProfileManyParams);
lwM2MTestClient.start(true);
awaitObserveReadAll(2, lwm2mDevice.getId().getId().toString());
awaitObserveReadAll(1, lwm2mDevice.getId().getId().toString());
awaitUpdateReg(3);
}

0
application/src/test/resources/lwm2m/3.xml → application/src/test/resources/lwm2m/3-1_2.xml

3
common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java

@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChainType;
import java.util.List;
import java.util.function.Predicate;
/**
* Created by ashvayka on 27.04.17.
@ -86,6 +87,8 @@ public interface RelationService {
ListenableFuture<List<EntityRelation>> findByRelationPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery);
ListenableFuture<List<EntityRelation>> findFilteredRelationsByPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery, Predicate<EntityRelation> relationFilter);
List<EntityRelation> findByRelationPathQuery(TenantId tenantId, EntityRelationPathQuery relationPathQuery);
// TODO: This method may be useful for some validations in the future

4
common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java

@ -65,6 +65,10 @@ public class CalculatedField extends BaseData<CalculatedFieldId> implements HasN
EntityType.DEVICE, EntityType.ASSET, EntityType.CUSTOMER, EntityType.TENANT
));
public static boolean isSupportedRefEntity(EntityId entity) {
return SUPPORTED_REFERENCED_ENTITIES.contains(entity.getEntityType());
}
private TenantId tenantId;
private EntityId entityId;

10
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java

@ -23,12 +23,13 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.Output;
import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.relation.RelationPathLevel;
import java.util.Map;
@Data
public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration {
public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration, ScheduledUpdateSupportedCalculatedFieldConfiguration {
@NotNull
private RelationPathLevel relation;
@ -40,11 +41,18 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A
private Output output;
private boolean useLatestTs;
private int scheduledUpdateInterval;
@Override
public CalculatedFieldType getType() {
return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
}
@Override
public boolean isScheduledUpdateEnabled() {
return true;
}
@Override
public void validate() {
relation.validate();

4
common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java

@ -186,8 +186,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
private long maxStateSizeInKBytes = 32;
@Schema(example = "2")
private long maxSingleValueArgumentSizeInKBytes = 2;
@Schema(example = "3600")
private long minAllowedDeduplicationIntervalInSecForCF = 3600;
@Schema(example = "60")
private long minAllowedDeduplicationIntervalInSecForCF = 60;
@Schema(example = "60")
private long minAggregationIntervalInSecForCF = 60;

10
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java

@ -224,8 +224,11 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
log.info("[{}] Closing old session: {}", registration.getEndpoint(), new UUID(oldSessionInfo.get().getSessionIdMSB(), oldSessionInfo.get().getSessionIdLSB()));
sessionManager.deregister(oldSessionInfo.get());
}
logService.log(lwM2MClient, LOG_LWM2M_INFO + ": Client registered with registration id: " + registration.getId() + " version: "
+ registration.getLwM2mVersion() + " and modes: " + registration.getQueueMode() + ", " + registration.getBindingMode());
String msgLogService = String.format("""
%s: Endpoint [%s] Client registered with registration id: [%s] LwM2mVersion: [%s], SupportedObjectIdVer [%s] QueueMode [%s], BindingMode %s
""", LOG_LWM2M_INFO, registration.getEndpoint(), registration.getId(), registration.getLwM2mVersion(), registration.getSupportedObject(), registration.getQueueMode(), registration.getBindingMode());
logService.log(lwM2MClient, msgLogService);
log.debug(msgLogService);
sessionManager.register(lwM2MClient.getSession());
this.initClientTelemetry(lwM2MClient);
this.initAttributes(lwM2MClient, true);
@ -244,7 +247,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
logService.log(lwM2MClient, LOG_LWM2M_WARN + ": Client registration failed due to invalid state: " + stateException.getState());
}
} catch (Throwable t) {
log.error("[{}] endpoint [{}] error Unable registration.", registration.getEndpoint(), t);
log.error("Endpoint [{}], Error Unable registration: [{}].", registration.getEndpoint(), t.getMessage(), t);
logService.log(lwM2MClient, LOG_LWM2M_WARN + ": Client registration failed due to: " + t.getMessage());
}
});
@ -290,7 +293,6 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
clientContext.unregister(client, registration);
SessionInfoProto sessionInfo = client.getSession();
if (sessionInfo != null) {
securityStore.remove(client.getEndpoint(), client.getRegistration().getId());
sessionManager.deregister(sessionInfo);
sessionStore.remove(registration.getEndpoint());
log.info("Client close session: [{}] unReg [{}] name [{}] profile ", registration.getId(), registration.getEndpoint(), sessionInfo.getDeviceType());

26
dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java

@ -71,6 +71,7 @@ import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.Predicate;
import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validatePositiveNumber;
@ -507,6 +508,11 @@ public class BaseRelationService implements RelationService {
@Override
public ListenableFuture<List<EntityRelation>> findByRelationPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery) {
return findFilteredRelationsByPathQueryAsync(tenantId, relationPathQuery, null);
}
@Override
public ListenableFuture<List<EntityRelation>> findFilteredRelationsByPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery, Predicate<EntityRelation> relationFilter) {
log.trace("Executing findByRelationPathQuery, tenantId [{}], relationPathQuery {}", tenantId, relationPathQuery);
validateId(tenantId, id -> "Invalid tenant id: " + id);
validate(relationPathQuery);
@ -518,10 +524,24 @@ public class BaseRelationService implements RelationService {
case FROM -> findByFromAndTypeAsync(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON);
case TO -> findByToAndTypeAsync(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON);
};
return Futures.transform(relationsFuture, entityRelations -> entityRelations.size() > limit ?
entityRelations.subList(0, limit) : entityRelations, MoreExecutors.directExecutor());
return Futures.transform(relationsFuture, entityRelations -> {
if (entityRelations == null || entityRelations.isEmpty()) {
return Collections.emptyList();
}
List<EntityRelation> relations = relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations;
return relations.size() > limit ? relations.subList(0, limit) : relations;
}, MoreExecutors.directExecutor());
}
return executor.submit(() -> relationDao.findByRelationPathQuery(tenantId, relationPathQuery, limit));
return executor.submit(() -> {
List<EntityRelation> entityRelations = relationDao.findByRelationPathQuery(tenantId, relationPathQuery, limit);
return relationFilter != null ? filterRelations(entityRelations, relationFilter) : entityRelations;
});
}
private List<EntityRelation> filterRelations(List<EntityRelation> entityRelations, Predicate<EntityRelation> relationFilter) {
return entityRelations.stream()
.filter(relationFilter)
.toList();
}
@Override

4
dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

@ -190,7 +190,7 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
userCredentialsDao.save(user.getTenantId(), userCredentials);
}
eventPublisher.publishEvent(SaveEntityEvent.builder()
.tenantId(tenantId == null ? TenantId.SYS_TENANT_ID : tenantId)
.tenantId(savedUser.getTenantId())
.entity(savedUser)
.oldEntity(oldUser)
.entityId(savedUser.getId())
@ -349,7 +349,7 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
eventPublisher.publishEvent(new UserCredentialsInvalidationEvent(userId));
countService.publishCountEntityEvictEvent(tenantId, EntityType.USER);
eventPublisher.publishEvent(DeleteEntityEvent.builder()
.tenantId(tenantId)
.tenantId(user.getTenantId())
.entityId(userId)
.entity(user)
.cause(cause)

1
dao/src/test/java/org/thingsboard/server/dao/service/OtaPackageServiceTest.java

@ -53,6 +53,7 @@ import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE;
public class OtaPackageServiceTest extends AbstractServiceTest {
public static final String TITLE = "My firmware";
public static final String TARGET_FW_VERSION = "fw.v.1.5.0-update";
private static final String FILE_NAME = "filename.txt";
private static final String VERSION = "v1.0";
private static final String CONTENT_TYPE = "text/plain";

6
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java

@ -449,8 +449,7 @@ public class CalculatedFieldTest extends AbstractContainerTest {
cf.setConfigurationVersion(1);
PropagationCalculatedFieldConfiguration cfg = new PropagationCalculatedFieldConfiguration();
cfg.setDirection(EntitySearchDirection.TO);
cfg.setRelationType(EntityRelation.CONTAINS_TYPE);
cfg.setRelation(new RelationPathLevel(EntitySearchDirection.TO, EntityRelation.CONTAINS_TYPE));
cfg.setApplyExpressionToResolvedArguments(true);
Argument arg = new Argument();
@ -535,8 +534,7 @@ public class CalculatedFieldTest extends AbstractContainerTest {
cf.setConfigurationVersion(1);
PropagationCalculatedFieldConfiguration cfg = new PropagationCalculatedFieldConfiguration();
cfg.setDirection(EntitySearchDirection.TO);
cfg.setRelationType(EntityRelation.CONTAINS_TYPE);
cfg.setRelation(new RelationPathLevel(EntitySearchDirection.TO, EntityRelation.CONTAINS_TYPE));
cfg.setApplyExpressionToResolvedArguments(false); // arguments-only mode
Argument arg = new Argument();

4
ui-ngx/src/app/modules/home/components/alarm-rules/alarm-rules-table-config.ts

@ -119,8 +119,8 @@ export class AlarmRulesTableConfig extends EntityTableConfig<any> {
this.columns.push(new EntityTableColumn<CalculatedFieldAlarmRule>('createRule', 'alarm-rule.severities', '67%',
entity => Object.keys(entity.configuration.createRules).map((severity) => this.translate.instant(alarmSeverityTranslations.get(severity as AlarmSeverity))).join(', '),
() => ({}), false));
this.columns.push(new EntityTableColumn<CalculatedFieldAlarmRule>('clearRule', 'alarm-rule.cleared', '60px',
entity => checkBoxCell(!!entity.configuration.clearRule), ()=> { return {padding: '0 14px'}}, false));
this.columns.push(new EntityTableColumn<CalculatedFieldAlarmRule>('clearRule', 'alarm-rule.cleared', '70px',
entity => checkBoxCell(!!entity.configuration.clearRule), ()=> { return {padding: 0, textAlign: 'center'}}, false));
this.cellActionDescriptors.push(
{

2
ui-ngx/src/app/modules/home/components/alarm-rules/cf-alarm-rule.component.html

@ -18,7 +18,7 @@
<div class="tb-form-panel no-border no-padding" [formGroup]="alarmRuleFormGroup">
<tb-cf-alarm-rule-condition formControlName="condition" [arguments]="arguments">
</tb-cf-alarm-rule-condition>
@if (!disabled || alarmRuleFormGroup.get('dashboardId').value) {
@if (!disabled || alarmRuleFormGroup.get('alarmDetails').value) {
<div class="tb-form-row space-between column-xs">
<div class="min-w-40 xs:min-w-fit" translate>
alarm-rule.alarm-rule-additional-info

3
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-complex-filter-predicate-dialog.component.ts

@ -15,7 +15,6 @@
///
import { Component, Inject } from '@angular/core';
import { ErrorStateMatcher } from '@angular/material/core';
import { MAT_DIALOG_DATA, MatDialogRef } from '@angular/material/dialog';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
@ -41,7 +40,7 @@ export interface AlarmRuleComplexFilterPredicateDialogData {
@Component({
selector: 'tb-alarm-rule-complex-filter-predicate-dialog',
templateUrl: './alarm-rule-complex-filter-predicate-dialog.component.html',
providers: [{provide: ErrorStateMatcher, useExisting: AlarmRuleComplexFilterPredicateDialogComponent}],
providers: [],
styleUrls: []
})

8
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-list.component.ts

@ -102,6 +102,14 @@ export class AlarmRuleFilterListComponent implements ControlValueAccessor, Valid
};
}
setDisabledState(isDisabled: boolean): void {
if (isDisabled) {
this.filterListFormGroup.disable({emitEvent: false});
} else {
this.filterListFormGroup.enable({emitEvent: false});
}
}
writeValue(filters: Array<AlarmRuleFilter>): void {
const keyFilterControls: Array<AbstractControl> = [];
if (filters) {

2
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-list.component.html

@ -46,7 +46,7 @@
[formControl]="predicateControl">
</tb-alarm-rule-filter-predicate>
<button mat-icon-button color="primary"
[class.!hidden]="disabled"
[disabled]="disabled"
type="button"
(click)="removePredicate($index)"
matTooltip="{{ 'filter.remove-filter' | translate }}"

2
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-list.component.ts

@ -107,7 +107,7 @@ export class AlarmRuleFilterPredicateListComponent implements ControlValueAccess
registerOnTouched(fn: any): void {
}
setDisabledState?(isDisabled: boolean): void {
setDisabledState(isDisabled: boolean): void {
this.disabled = isDisabled;
if (this.disabled) {
this.filterListFormGroup.disable({emitEvent: false});

10
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate-value.component.ts

@ -107,6 +107,16 @@ export class AlarmRuleFilterPredicateValueComponent implements ControlValueAcces
).subscribe(value => this.updateValueModeValidators(value))
}
setDisabledState(isDisabled: boolean): void {
if (isDisabled) {
this.filterPredicateValueFormGroup.disable({emitEvent: false});
this.dynamicModeControl.disable({emitEvent: false});
} else {
this.filterPredicateValueFormGroup.enable({emitEvent: false});
this.dynamicModeControl.enable({emitEvent: false});
}
}
private updateValueModeValidators(isDynamicMode: boolean): void {
if (isDynamicMode) {
this.filterPredicateValueFormGroup.get('staticValue').disable({emitEvent: false});

14
ui-ngx/src/app/modules/home/components/alarm-rules/filter/alarm-rule-filter-predicate.component.ts

@ -14,7 +14,7 @@
/// limitations under the License.
///
import { Component, DestroyRef, forwardRef, Input } from '@angular/core';
import { booleanAttribute, Component, DestroyRef, forwardRef, Input } from '@angular/core';
import {
ControlValueAccessor,
FormBuilder,
@ -61,6 +61,9 @@ import { CalculatedFieldArgument } from "@shared/models/calculated-field.models"
})
export class AlarmRuleFilterPredicateComponent implements ControlValueAccessor, Validator {
@Input({ transform: booleanAttribute })
disabled: boolean;
@Input()
valueType: EntityKeyValueType;
@ -115,6 +118,15 @@ export class AlarmRuleFilterPredicateComponent implements ControlValueAccessor,
};
}
setDisabledState(isDisabled: boolean): void {
this.disabled = isDisabled;
if (isDisabled) {
this.filterPredicateFormGroup.disable({emitEvent: false});
} else {
this.filterPredicateFormGroup.enable({emitEvent: false});
}
}
writeValue(predicate: AlarmRuleFilterPredicate): void {
this.type = predicate.type;
this.filterPredicateFormGroup.patchValue(predicate, {emitEvent: false});

1
ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html

@ -70,6 +70,7 @@
<tb-geofencing-configuration formControlName="configuration"
[entityId]="data.entityId"
[entityName]="data.entityName"
[ownerId]="data.ownerId"
[tenantId]="data.tenantId"
></tb-geofencing-configuration>
}

12
ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/calculated-field-geofencing-zone-groups-panel.component.ts

@ -69,6 +69,7 @@ export class CalculatedFieldGeofencingZoneGroupsPanelComponent implements OnInit
@Input() entityId: EntityId;
@Input() tenantId: string;
@Input() entityName: string;
@Input() ownerId: EntityId;
@Input() usedNames: string[];
@ViewChild('entityAutocomplete') entityAutocomplete: EntityAutocompleteComponent;
@ -204,6 +205,11 @@ export class CalculatedFieldGeofencingZoneGroupsPanelComponent implements OnInit
case ArgumentEntityType.Current:
delete value.refEntityId;
break;
case ArgumentEntityType.Owner:
delete value.refEntityId;
value.refDynamicSourceConfiguration ||= { type: ArgumentEntityType.Owner };
value.refDynamicSourceConfiguration.type = ArgumentEntityType.Owner;
break;
case ArgumentEntityType.RelationQuery:
delete value.refEntityId;
value.refDynamicSourceConfiguration.type = ArgumentEntityType.RelationQuery;
@ -228,6 +234,12 @@ export class CalculatedFieldGeofencingZoneGroupsPanelComponent implements OnInit
case ArgumentEntityType.RelationQuery:
entityFilter = this.currentEntityFilter;
break;
case ArgumentEntityType.Owner:
entityFilter = {
type: AliasFilterType.singleEntity,
singleEntity: this.ownerId
};
break;
case ArgumentEntityType.Tenant:
entityFilter = {
type: AliasFilterType.singleEntity,

2
ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/calculated-field-geofencing-zone-groups-table.component.ts

@ -81,6 +81,7 @@ export class CalculatedFieldGeofencingZoneGroupsTableComponent implements Contro
@Input({required: true}) entityId: EntityId;
@Input({required: true}) tenantId: string;
@Input({required: true}) entityName: string;
@Input({required: true}) ownerId: EntityId;
@ViewChild(MatSort, { static: true }) sort: MatSort;
@ -159,6 +160,7 @@ export class CalculatedFieldGeofencingZoneGroupsTableComponent implements Contro
buttonTitle: isExists ? 'action.apply' : 'action.add',
tenantId: this.tenantId,
entityName: this.entityName,
ownerId: this.ownerId,
usedNames: this.zoneGroupsFormArray.value.map(({ name }) => name).filter(name => name !== zone.name),
};
this.popoverComponent = this.popoverService.displayPopover({

1
ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/geofencing-configuration.component.html

@ -43,6 +43,7 @@
<tb-calculated-field-geofencing-zone-groups-table formControlName="zoneGroups"
[entityId]="entityId"
[tenantId]="tenantId"
[ownerId]="ownerId"
[entityName]="entityName"/>
<div class="tb-form-row space-between flex-1 columns-xs" [class.!hidden]="!isRelatedEntity">
<mat-slide-toggle class="mat-slide" formControlName="scheduledUpdateEnabled">

3
ui-ngx/src/app/modules/home/components/calculated-fields/components/geofencing-configuration/geofencing-configuration.component.ts

@ -69,6 +69,9 @@ export class GeofencingConfigurationComponent implements ControlValueAccessor, V
@Input({required: true})
entityName: string;
@Input({required: true})
ownerId: EntityId;
readonly minAllowedScheduledUpdateIntervalInSecForCF = getCurrentAuthState(this.store).minAllowedScheduledUpdateIntervalInSecForCF;
readonly DataKeyType = DataKeyType;

2
ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts

@ -105,6 +105,7 @@ export class RelatedEntitiesAggregationComponentComponent implements ControlValu
map(argumentsObj => getCalculatedFieldArgumentsHighlights(argumentsObj))
);
private readonly minAllowedScheduledUpdateIntervalInSecForCF = getCurrentAuthState(this.store).minAllowedScheduledUpdateIntervalInSecForCF;
private propagateChange: (config: CalculatedFieldRelatedAggregationConfiguration) => void = () => { };
constructor(private fb: FormBuilder,
@ -149,6 +150,7 @@ export class RelatedEntitiesAggregationComponentComponent implements ControlValu
private updatedModel(value: CalculatedFieldRelatedAggregationConfiguration): void {
value.type = CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
value.scheduledUpdateInterval = this.minAllowedScheduledUpdateIntervalInSecForCF;
this.propagateChange(value);
}
}

2
ui-ngx/src/app/modules/home/pages/asset/asset-tabs.component.html

@ -37,7 +37,7 @@
</mat-tab>
<mat-tab #alarmRules="matTab"
label="{{'alarm-rule.alarm-rules' | translate }}">
<tb-alarm-rules-table [active]="alarmRules.isActive" [entityId]="entity.id" [entityName]="entity.name"></tb-alarm-rules-table>
<tb-alarm-rules-table [active]="alarmRules.isActive" [entityId]="entity.id" [entityName]="entity.name" [ownerId]="entity.tenantId"></tb-alarm-rules-table>
</mat-tab>
}
<mat-tab label="{{ 'alarm.alarms' | translate }}" #alarmsTab="matTab">

12
ui-ngx/src/app/modules/home/pages/device-profile/device-profile-tabs.component.html

@ -50,14 +50,14 @@
<mat-tab #alarmRules="matTab" label="{{'alarm-rule.alarm-rules-tab' | translate }}">
<section class="flex flex-col h-full">
@if (hasOldRules && !isEdit) {
<div class="flex flex-row items-center justify-end" style="padding: 16px 16px 0;">
<div class="flex flex-row items-center justify-start" style="padding: 16px 16px 0;">
<tb-toggle-select [(ngModel)]="alarmRulesOldVersion" appearance="fill">
<tb-toggle-option [value]="false">{{ 'alarm-rule.alarm-rules-new' | translate }}</tb-toggle-option>
<tb-toggle-option [value]="false">{{ 'alarm-rule.alarm-rules-actual' | translate }}</tb-toggle-option>
<tb-toggle-option [value]="true">{{ 'alarm-rule.alarm-rules-old' | translate }}</tb-toggle-option>
</tb-toggle-select>
</div>
}
@if (alarmRulesOldVersion) {
@if (alarmRulesOldVersion || isEdit) {
<div class="mat-padding" [formGroup]="detailsForm" *ngIf="alarmRules.isActive">
<div formGroupName="profileData">
<tb-device-profile-alarms formControlName="alarms" [deviceProfileId]="entity.id"></tb-device-profile-alarms>
@ -65,7 +65,11 @@
</div>
} @else {
<div class="relative h-full">
<tb-alarm-rules-table [active]="alarmRules.isActive" [entityId]="entity.id" [entityName]="entity.name" [ownerId]="entity.tenantId"></tb-alarm-rules-table>
<tb-alarm-rules-table [active]="alarmRules.isActive"
[entityId]="entity.id"
[entityName]="entity.name"
[ownerId]="entity.tenantId">
</tb-alarm-rules-table>
</div>
}
</section>

2
ui-ngx/src/app/modules/home/pages/device-profile/device-profile-tabs.component.ts

@ -61,7 +61,7 @@ export class DeviceProfileTabsComponent extends EntityTabsComponent<DeviceProfil
protected setEntity(entity: DeviceProfile) {
this.isTransportTypeChanged = false;
this.hasOldRules = !!entity?.profileData?.alarms?.length;
this.alarmRulesOldVersion = this.isEdit || false;
this.alarmRulesOldVersion = false;
super.setEntity(entity);
}

4
ui-ngx/src/app/modules/home/pages/device/device-tabs.component.html

@ -33,10 +33,10 @@
</mat-tab>
@if (authUser.authority === authorities.TENANT_ADMIN) {
<mat-tab label="{{ 'entity.type-calculated-fields' | translate }}" #calculatedFieldsTab="matTab">
<tb-calculated-fields-table [active]="calculatedFieldsTab.isActive" [entityId]="entity.id"/>
<tb-calculated-fields-table [active]="calculatedFieldsTab.isActive" [entityId]="entity.id" [ownerId]="entity.tenantId"/>
</mat-tab>
<mat-tab #calculatedFieldsAlarmRules="matTab" label="{{'alarm-rule.alarm-rules-tab' | translate }}">
<tb-alarm-rules-table [active]="calculatedFieldsAlarmRules.isActive" [entityId]="entity.id" [entityName]="entity.name"></tb-alarm-rules-table>
<tb-alarm-rules-table [active]="calculatedFieldsAlarmRules.isActive" [entityId]="entity.id" [entityName]="entity.name" [ownerId]="entity.tenantId"></tb-alarm-rules-table>
</mat-tab>
}
<mat-tab label="{{ 'alarm.alarms' | translate }}" #alarmsTab="matTab">

1
ui-ngx/src/app/shared/models/calculated-field.models.ts

@ -150,6 +150,7 @@ export interface CalculatedFieldRelatedAggregationConfiguration {
arguments: Record<string, CalculatedFieldArgument>;
metrics: Record<string, CalculatedFieldAggMetric>;
deduplicationIntervalInSec: number;
scheduledUpdateInterval?: number;
useLatestTs: boolean;
output: CalculatedFieldOutput & { decimalsByDefault?: number; };
}

2
ui-ngx/src/app/shared/models/tenant.model.ts

@ -176,7 +176,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan
maxArgumentsPerCF: 10,
maxDataPointsPerRollingArg: 1000,
maxRelationLevelPerCfArgument: 10,
minAllowedDeduplicationIntervalInSecForCF: 3600,
minAllowedDeduplicationIntervalInSecForCF: 60,
minAggregationIntervalInSecForCF: 60,
maxRelatedEntitiesToReturnPerCfArgument: 100,
minAllowedScheduledUpdateIntervalInSecForCF: 0,

4
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -1288,9 +1288,9 @@
"alarm-rule": "Alarm rule",
"alarm-rules": "Alarm rules",
"alarm-rules-old": "Old",
"alarm-rules-new": "New",
"alarm-rules-actual": "Actual",
"severities": "Severities",
"cleared": "Cleared",
"cleared": "Clear condition",
"delete-title": "Are you sure you want to delete the alarm rule '{{title}}'?",
"delete-text": "Be careful, after the confirmation the alarm rule and all related data will become unrecoverable.",
"delete-multiple-title": "Are you sure you want to delete { count, plural, =1 {1 alarm rule} other {# alarm rules} }?",

Loading…
Cancel
Save