Browse Source

Merge branch 'master' into sparkplug-any-name-node-unique-device-names

# Conflicts:
#	application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/AbstractMqttV5ClientSparkplugTest.java
pull/15041/head
nickAS21 7 months ago
parent
commit
a25c5ca66b
  1. 1
      .gitignore
  2. 39
      TEST_FAST.md
  3. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java
  4. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  7. 35
      application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextFactory.java
  8. 6
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  9. 12
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractTbRuleEngineSubmitStrategy.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/queue/processing/BatchTbRuleEngineSubmitStrategy.java
  11. 2
      application/src/main/java/org/thingsboard/server/service/queue/processing/BurstTbRuleEngineSubmitStrategy.java
  12. 14
      application/src/main/java/org/thingsboard/server/service/queue/processing/IdMsgPair.java
  13. 10
      application/src/main/java/org/thingsboard/server/service/queue/processing/SequentialByEntityIdTbRuleEngineSubmitStrategy.java
  14. 6
      application/src/main/java/org/thingsboard/server/service/queue/processing/SequentialTbRuleEngineSubmitStrategy.java
  15. 8
      application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java
  16. 10
      application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationDetails.java
  17. 7
      application/src/main/resources/thingsboard.yml
  18. 1
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  19. 17
      application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java
  20. 41
      application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java
  21. 2
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesAggregationCalculatedFieldStateTest.java
  22. 260
      application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java
  23. 281
      application/src/test/java/org/thingsboard/server/service/queue/ruleengine/RuleEngineConsumerLoopTest.java
  24. 2
      application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java
  25. 2
      application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineStrategyTest.java
  26. 1
      application/src/test/java/org/thingsboard/server/service/sms/DefaultSmsServiceTest.java
  27. 1
      application/src/test/java/org/thingsboard/server/service/stats/DevicesStatisticsTest.java
  28. 52
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/FwLwM2MDevice.java
  29. 60
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SwLwM2MDevice.java
  30. 13
      common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java
  31. 8
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java
  32. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java
  33. 5
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java
  34. 2
      common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java
  35. 4
      common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfigurationTest.java
  36. 14
      common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java
  37. 8
      common/edge-api/src/main/proto/edge.proto
  38. 4
      common/proto/src/main/java/org/thingsboard/server/common/adaptor/ProtoConverter.java
  39. 20
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java
  40. 14
      dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java
  41. 4
      dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
  42. 232
      dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java
  43. 56
      dao/src/test/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidatorTest.java
  44. 2
      edqs/src/main/resources/edqs.yml
  45. 4
      msa/vc-executor/src/main/resources/tb-vc-executor.yml
  46. 43
      rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java
  47. 31
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeTest.java
  48. 4
      transport/coap/src/main/resources/tb-coap-transport.yml
  49. 4
      transport/http/src/main/resources/tb-http-transport.yml
  50. 4
      transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml
  51. 4
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml
  52. 4
      transport/snmp/src/main/resources/tb-snmp-transport.yml
  53. 31
      ui-ngx/patches/ngx-hm-carousel+19.0.0.patch
  54. 4
      ui-ngx/src/app/modules/home/components/alarm-rules/alarm-rules-table-config.ts
  55. 1
      ui-ngx/src/app/modules/home/components/api-key/api-keys-table-config.ts
  56. 2
      ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html
  57. 12
      ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.ts
  58. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/calculated-fields-table-config.ts
  59. 104
      ui-ngx/src/app/modules/home/components/entity/entities-table.component.html
  60. 2
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.scss
  61. 18
      ui-ngx/src/app/modules/home/components/vc/entity-versions-table.component.html
  62. 11
      ui-ngx/src/app/modules/home/components/widget/lib/home-page/recent-dashboards-widget.component.html
  63. 45
      ui-ngx/src/app/modules/home/components/widget/lib/home-page/recent-dashboards-widget.component.ts
  64. 40
      ui-ngx/src/app/modules/home/components/widget/lib/rpc/power-button-widget.models.ts
  65. 2
      ui-ngx/src/app/modules/home/components/widget/lib/settings/common/image-cards-select.component.scss
  66. 5
      ui-ngx/src/app/modules/home/components/widget/lib/settings/common/map/marker-image-settings.component.ts
  67. 5
      ui-ngx/src/app/modules/home/components/widget/lib/settings/common/map/shape-fill-image-settings.component.ts
  68. 2
      ui-ngx/src/app/modules/home/components/widget/lib/settings/widget-settings.scss
  69. 6
      ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts
  70. 2
      ui-ngx/src/app/modules/home/models/entity/entities-table-config.models.ts
  71. 2
      ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html
  72. 1
      ui-ngx/src/app/modules/home/pages/ai-model/ai-model-table-config.resolve.ts
  73. 1
      ui-ngx/src/app/modules/home/pages/mobile/applications/mobile-app-table-config.resolver.ts
  74. 1
      ui-ngx/src/app/modules/home/pages/mobile/bundes/mobile-bundle-table-config.resolve.ts
  75. 2
      ui-ngx/src/app/modules/home/pages/mobile/common/editor-panel.component.ts
  76. 1
      ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-table-config.resolver.ts
  77. 1
      ui-ngx/src/app/modules/home/pages/notification/rule/rule-table-config.resolver.ts
  78. 2
      ui-ngx/src/app/modules/home/pages/notification/template/configuration/notification-template-configuration.component.scss
  79. 1
      ui-ngx/src/app/modules/home/pages/notification/template/template-table-config.resolver.ts
  80. 2
      ui-ngx/src/app/modules/login/pages/login/login.component.ts
  81. 7
      ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts
  82. 45
      ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts
  83. 43
      ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay.ts
  84. 6
      ui-ngx/src/app/shared/components/image/image-gallery.component.html
  85. 1
      ui-ngx/src/app/shared/models/entity-type.models.ts
  86. 2
      ui-ngx/src/assets/dashboard/tenant_admin_home_page.json
  87. 1
      ui-ngx/src/assets/locale/locale.constant-en_US.json

1
.gitignore

@ -37,3 +37,4 @@ rebuild-docker.sh
*/.run/**
.run/**
.run
.claude/

39
TEST_FAST.md

@ -10,22 +10,43 @@ mvn clean install -T6 -DskipTests
mvn test -pl='!application,!dao,!ui-ngx,!msa/js-executor,!msa/web-ui' -T4
mvn test -pl dao -Dparallel=packages -DforkCount=4
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.controller.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.edge.**' -DforkCount=4 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.service.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.mqtt.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.coap.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='org.thingsboard.server.transport.lwm2m.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dsurefire.excludes='**/nosql/*Test.java' -Dtest='**/*TestSuite.java' -DforkCount=4 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.controller.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.edge.**' -DforkCount=4 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.service.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.transport.mqtt.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.transport.coap.**' -DforkCount=6 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='!**/nosql/**,org.thingsboard.server.transport.lwm2m.**' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
mvn test -pl application -Dtest='**/*TestSuite.java' -DforkCount=4 -Dparallel=classes -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
#the rest of application tests
mvn test -pl application -Dtest='
!**/nosql/*Test.java,
!**/nosql/**,
!org.thingsboard.server.controller.**,
!org.thingsboard.server.edge.**,
!org.thingsboard.server.service.**,
!org.thingsboard.server.transport.mqtt.**,
!org.thingsboard.server.transport.coap.**,
!org.thingsboard.server.transport.lwm2m.**
!org.thingsboard.server.transport.lwm2m.**,
!**/*TestSuite.java
' -DforkCount=6 -Dparallel=packages -Dsurefire.rerunFailingTestsCount=2 -Dsurefire.failOnFlakeCount=5
```
## Testcontainers compatibility with the Docker API workaround
In case your tests failed to run testcontainers due to unsupported Docker API version
:coffee: testcontainers (Docker API 1.32) + :whale: docker 29 (min API 1.44) workaround
Add to /etc/docker/daemon.json and restart docker
```json
{
"min-api-version": "1.32"
}
```
Same works on Mac, except `daemon.json` are located in another folder and required to be edited from Docker Desktop UI.
Tip: If your testcontainer are struggling to find any Docker. You can try to remove the testcontainers property file. It will be recreated on the next testcontainers run.
```bash
rm ~/.testcontainers.properties
```

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java

@ -147,15 +147,15 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa
private void processMsgs(List<TbProtoQueueMsg<ToCalculatedFieldMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumer, Object consumerKey, QueueConfig config) throws Exception {
List<IdMsgPair<ToCalculatedFieldMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList();
ConcurrentMap<UUID, TbProtoQueueMsg<ToCalculatedFieldMsg>> pendingMap = orderedMsgList.stream().collect(
Collectors.toConcurrentMap(IdMsgPair::getUuid, IdMsgPair::getMsg));
Collectors.toConcurrentMap(IdMsgPair::uuid, IdMsgPair::msg));
CountDownLatch processingTimeoutLatch = new CountDownLatch(1);
TbPackProcessingContext<TbProtoQueueMsg<ToCalculatedFieldMsg>> ctx = new TbPackProcessingContext<>(
processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>());
PendingMsgHolder<ToCalculatedFieldMsg> pendingMsgHolder = new PendingMsgHolder<>();
Future<?> packSubmitFuture = consumersExecutor.submit(() -> {
orderedMsgList.forEach((element) -> {
UUID id = element.getUuid();
TbProtoQueueMsg<ToCalculatedFieldMsg> msg = element.getMsg();
UUID id = element.uuid();
TbProtoQueueMsg<ToCalculatedFieldMsg> msg = element.msg();
log.trace("[{}] Creating main callback for message: {}", id, msg.getValue());
TbCallback callback = new TbPackCallback<>(id, ctx);
try {

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -260,15 +260,15 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private void processMsgs(List<TbProtoQueueMsg<ToCoreMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToCoreMsg>> consumer, Object consumerKey, QueueConfig config) throws Exception {
List<IdMsgPair<ToCoreMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList();
ConcurrentMap<UUID, TbProtoQueueMsg<ToCoreMsg>> pendingMap = orderedMsgList.stream().collect(
Collectors.toConcurrentMap(IdMsgPair::getUuid, IdMsgPair::getMsg));
Collectors.toConcurrentMap(IdMsgPair::uuid, IdMsgPair::msg));
CountDownLatch processingTimeoutLatch = new CountDownLatch(1);
TbPackProcessingContext<TbProtoQueueMsg<ToCoreMsg>> ctx = new TbPackProcessingContext<>(
processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>());
PendingMsgHolder<ToCoreMsg> pendingMsgHolder = new PendingMsgHolder<>();
Future<?> packSubmitFuture = consumersExecutor.submit(() -> {
orderedMsgList.forEach((element) -> {
UUID id = element.getUuid();
TbProtoQueueMsg<ToCoreMsg> msg = element.getMsg();
UUID id = element.uuid();
TbProtoQueueMsg<ToCoreMsg> msg = element.msg();
log.trace("[{}] Creating main callback for message: {}", id, msg.getValue());
TbCallback callback = new TbPackCallback<>(id, ctx);
try {

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java

@ -128,15 +128,15 @@ public class DefaultTbEdgeConsumerService extends AbstractConsumerService<ToEdge
private void processMsgs(List<TbProtoQueueMsg<ToEdgeMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToEdgeMsg>> consumer, Object consumerKey, QueueConfig edgeQueueConfig) throws InterruptedException {
List<IdMsgPair<ToEdgeMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList();
ConcurrentMap<UUID, TbProtoQueueMsg<ToEdgeMsg>> pendingMap = orderedMsgList.stream().collect(
Collectors.toConcurrentMap(IdMsgPair::getUuid, IdMsgPair::getMsg));
Collectors.toConcurrentMap(IdMsgPair::uuid, IdMsgPair::msg));
CountDownLatch processingTimeoutLatch = new CountDownLatch(1);
TbPackProcessingContext<TbProtoQueueMsg<ToEdgeMsg>> ctx = new TbPackProcessingContext<>(
processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>());
PendingMsgHolder<ToEdgeMsg> pendingMsgHolder = new PendingMsgHolder<>();
Future<?> submitFuture = consumersExecutor.submit(() -> {
orderedMsgList.forEach((element) -> {
UUID id = element.getUuid();
TbProtoQueueMsg<ToEdgeMsg> msg = element.getMsg();
UUID id = element.uuid();
TbProtoQueueMsg<ToEdgeMsg> msg = element.msg();
TbCallback callback = new TbPackCallback<>(id, ctx);
try {
ToEdgeMsg toEdgeMsg = msg.getValue();

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -70,6 +70,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedCo
private final TbRuleEngineConsumerContext ctx;
private final QueueService queueService;
private final TbRuleEngineDeviceRpcService tbDeviceRpcService;
private final TbMsgPackProcessingContextFactory packProcessingContextFactory;
private final ConcurrentMap<QueueKey, TbRuleEngineQueueConsumerManager> consumers = new ConcurrentHashMap<>();
@ -85,11 +86,13 @@ public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedCo
PartitionService partitionService,
ApplicationEventPublisher eventPublisher,
JwtSettingsService jwtSettingsService,
CalculatedFieldCache calculatedFieldCache) {
CalculatedFieldCache calculatedFieldCache,
TbMsgPackProcessingContextFactory packProcessingContextFactory) {
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, tbResourceDataCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService);
this.ctx = ctx;
this.tbDeviceRpcService = tbDeviceRpcService;
this.queueService = queueService;
this.packProcessingContextFactory = packProcessingContextFactory;
}
@Override
@ -255,6 +258,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedCo
.consumerExecutor(consumersExecutor)
.scheduler(scheduler)
.taskExecutor(mgmtExecutor)
.packProcessingContextFactory(packProcessingContextFactory)
.build();
consumers.put(queueKey, consumer);
consumer.init(queue);

35
application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextFactory.java

@ -0,0 +1,35 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.queue;
import org.springframework.stereotype.Component;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;
public interface TbMsgPackProcessingContextFactory {
TbMsgPackProcessingContext create(String queueName, TbRuleEngineSubmitStrategy submitStrategy, boolean skipTimeouts);
@Component
class DefaultTbMsgPackProcessingContextFactory implements TbMsgPackProcessingContextFactory {
@Override
public TbMsgPackProcessingContext create(String queueName, TbRuleEngineSubmitStrategy submitStrategy, boolean skipTimeouts) {
return new TbMsgPackProcessingContext(queueName, submitStrategy, skipTimeouts);
}
}
}

6
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -134,13 +134,13 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
protected void processNotifications(List<TbProtoQueueMsg<N>> msgs, TbQueueConsumer<TbProtoQueueMsg<N>> consumer) throws Exception {
List<IdMsgPair<N>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList();
ConcurrentMap<UUID, TbProtoQueueMsg<N>> pendingMap = orderedMsgList.stream().collect(
Collectors.toConcurrentMap(IdMsgPair::getUuid, IdMsgPair::getMsg));
Collectors.toConcurrentMap(IdMsgPair::uuid, IdMsgPair::msg));
CountDownLatch processingTimeoutLatch = new CountDownLatch(1);
TbPackProcessingContext<TbProtoQueueMsg<N>> ctx = new TbPackProcessingContext<>(
processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>());
orderedMsgList.forEach(element -> {
UUID id = element.getUuid();
TbProtoQueueMsg<N> msg = element.getMsg();
UUID id = element.uuid();
TbProtoQueueMsg<N> msg = element.msg();
log.trace("[{}] Creating notification callback for message: {}", id, msg.getValue());
TbCallback callback = new TbPackCallback<>(id, ctx);
try {

12
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractTbRuleEngineSubmitStrategy.java

@ -44,21 +44,21 @@ public abstract class AbstractTbRuleEngineSubmitStrategy implements TbRuleEngine
@Override
public ConcurrentMap<UUID, TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg>> getPendingMap() {
return orderedMsgList.stream().collect(Collectors.toConcurrentMap(pair -> pair.uuid, pair -> pair.msg));
return orderedMsgList.stream().collect(Collectors.toConcurrentMap(pair -> pair.uuid(), pair -> pair.msg()));
}
@Override
public void update(ConcurrentMap<UUID, TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg>> reprocessMap) {
List<IdMsgPair<TransportProtos.ToRuleEngineMsg>> newOrderedMsgList = new ArrayList<>(reprocessMap.size());
for (IdMsgPair<TransportProtos.ToRuleEngineMsg> pair : orderedMsgList) {
if (reprocessMap.containsKey(pair.uuid)) {
if (StringUtils.isNotEmpty(pair.getMsg().getValue().getFailureMessage())) {
var toRuleEngineMsg = TransportProtos.ToRuleEngineMsg.newBuilder(pair.getMsg().getValue())
if (reprocessMap.containsKey(pair.uuid())) {
if (StringUtils.isNotEmpty(pair.msg().getValue().getFailureMessage())) {
var toRuleEngineMsg = TransportProtos.ToRuleEngineMsg.newBuilder(pair.msg().getValue())
.clearFailureMessage()
.clearRelationTypes()
.build();
var newMsg = new TbProtoQueueMsg<>(pair.getMsg().getKey(), toRuleEngineMsg, pair.getMsg().getHeaders());
newOrderedMsgList.add(new IdMsgPair<>(pair.getUuid(), newMsg));
var newMsg = new TbProtoQueueMsg<>(pair.msg().getKey(), toRuleEngineMsg, pair.msg().getHeaders());
newOrderedMsgList.add(new IdMsgPair<>(pair.uuid(), newMsg));
} else {
newOrderedMsgList.add(pair);
}

2
application/src/main/java/org/thingsboard/server/service/queue/processing/BatchTbRuleEngineSubmitStrategy.java

@ -73,7 +73,7 @@ public class BatchTbRuleEngineSubmitStrategy extends AbstractTbRuleEngineSubmitS
pendingPack.clear();
for (int i = startIdx; i < endIdx; i++) {
IdMsgPair<TransportProtos.ToRuleEngineMsg> pair = orderedMsgList.get(i);
pendingPack.put(pair.uuid, pair.msg);
pendingPack.put(pair.uuid(), pair.msg());
}
tmpPack = new LinkedHashMap<>(pendingPack);
}

2
application/src/main/java/org/thingsboard/server/service/queue/processing/BurstTbRuleEngineSubmitStrategy.java

@ -34,7 +34,7 @@ public class BurstTbRuleEngineSubmitStrategy extends AbstractTbRuleEngineSubmitS
if (log.isDebugEnabled()) {
log.debug("[{}] submitting [{}] messages to rule engine", queueName, orderedMsgList.size());
}
orderedMsgList.forEach(pair -> msgConsumer.accept(pair.uuid, pair.msg));
orderedMsgList.forEach(pair -> msgConsumer.accept(pair.uuid(), pair.msg()));
}
@Override

14
application/src/main/java/org/thingsboard/server/service/queue/processing/IdMsgPair.java

@ -15,19 +15,9 @@
*/
package org.thingsboard.server.service.queue.processing;
import lombok.Getter;
import com.google.protobuf.GeneratedMessageV3;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import java.util.UUID;
public class IdMsgPair<T extends com.google.protobuf.GeneratedMessageV3> {
@Getter
final UUID uuid;
@Getter
final TbProtoQueueMsg<T> msg;
public IdMsgPair(UUID uuid, TbProtoQueueMsg<T> msg) {
this.uuid = uuid;
this.msg = msg;
}
}
public record IdMsgPair<T extends GeneratedMessageV3>(UUID uuid, TbProtoQueueMsg<T> msg) {}

10
application/src/main/java/org/thingsboard/server/service/queue/processing/SequentialByEntityIdTbRuleEngineSubmitStrategy.java

@ -51,7 +51,7 @@ public abstract class SequentialByEntityIdTbRuleEngineSubmitStrategy extends Abs
entityIdToListMap.forEach((entityId, queue) -> {
IdMsgPair<TransportProtos.ToRuleEngineMsg> msg = queue.peek();
if (msg != null) {
msgConsumer.accept(msg.uuid, msg.msg);
msgConsumer.accept(msg.uuid(), msg.msg());
}
});
}
@ -71,13 +71,13 @@ public abstract class SequentialByEntityIdTbRuleEngineSubmitStrategy extends Abs
IdMsgPair<TransportProtos.ToRuleEngineMsg> next = null;
synchronized (queue) {
IdMsgPair<TransportProtos.ToRuleEngineMsg> expected = queue.peek();
if (expected != null && expected.uuid.equals(id)) {
if (expected != null && expected.uuid().equals(id)) {
queue.poll();
next = queue.peek();
}
}
if (next != null) {
msgConsumer.accept(next.uuid, next.msg);
msgConsumer.accept(next.uuid(), next.msg());
}
}
}
@ -87,9 +87,9 @@ public abstract class SequentialByEntityIdTbRuleEngineSubmitStrategy extends Abs
msgToEntityIdMap.clear();
entityIdToListMap.clear();
for (IdMsgPair<TransportProtos.ToRuleEngineMsg> pair : orderedMsgList) {
EntityId entityId = getEntityId(pair.msg.getValue());
EntityId entityId = getEntityId(pair.msg().getValue());
if (entityId != null) {
msgToEntityIdMap.put(pair.uuid, entityId);
msgToEntityIdMap.put(pair.uuid(), entityId);
entityIdToListMap.computeIfAbsent(entityId, id -> new LinkedList<>()).add(pair);
}
}

6
application/src/main/java/org/thingsboard/server/service/queue/processing/SequentialTbRuleEngineSubmitStrategy.java

@ -60,11 +60,11 @@ public class SequentialTbRuleEngineSubmitStrategy extends AbstractTbRuleEngineSu
int idx = msgIdx.get();
if (idx < listSize) {
IdMsgPair<TransportProtos.ToRuleEngineMsg> pair = orderedMsgList.get(idx);
expectedMsgId = pair.uuid;
expectedMsgId = pair.uuid();
if (log.isDebugEnabled()) {
log.debug("[{}] submitting [{}] message to rule engine", queueName, pair.msg);
log.debug("[{}] submitting [{}] message to rule engine", queueName, pair.msg());
}
msgConsumer.accept(pair.uuid, pair.msg);
msgConsumer.accept(pair.uuid(), pair.msg());
}
}

8
application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java

@ -42,6 +42,7 @@ import org.thingsboard.server.queue.common.consumer.TbQueueConsumerTask.Consumer
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.service.queue.TbMsgPackCallback;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContext;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContextFactory;
import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingDecision;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingResult;
@ -67,13 +68,15 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager<T
private final TbRuleEngineConsumerContext ctx;
private final TbRuleEngineConsumerStats stats;
private final TbMsgPackProcessingContextFactory packProcessingContextFactory;
@Builder(builderMethodName = "create") // not to conflict with super.builder()
public TbRuleEngineQueueConsumerManager(TbRuleEngineConsumerContext ctx,
QueueKey queueKey,
ExecutorService consumerExecutor,
ScheduledExecutorService scheduler,
ExecutorService taskExecutor) {
ExecutorService taskExecutor,
TbMsgPackProcessingContextFactory packProcessingContextFactory) {
super(queueKey, null, null,
(queueConfig, tpi) -> {
Integer partitionId = tpi != null ? tpi.getPartition().orElse(-1) : null;
@ -82,6 +85,7 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager<T
consumerExecutor, scheduler, taskExecutor, null);
this.ctx = ctx;
this.stats = new TbRuleEngineConsumerStats(queueKey, ctx.getStatsFactory());
this.packProcessingContextFactory = packProcessingContextFactory;
}
public void delete(boolean drainQueue) {
@ -134,7 +138,7 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager<T
TbRuleEngineProcessingStrategy ackStrategy = getProcessingStrategy(queue);
submitStrategy.init(msgs);
while (!stopped && !consumer.isStopped()) {
TbMsgPackProcessingContext packCtx = new TbMsgPackProcessingContext(queue.getName(), submitStrategy, ackStrategy.isSkipTimeoutMsgs());
TbMsgPackProcessingContext packCtx = packProcessingContextFactory.create(queue.getName(), submitStrategy, ackStrategy.isSkipTimeoutMsgs());
submitStrategy.submitAttempt((id, msg) -> submitMessage(packCtx, id, msg));
final boolean timeout = !packCtx.await(queue.getPackProcessingTimeout(), TimeUnit.MILLISECONDS);

10
application/src/main/java/org/thingsboard/server/service/security/auth/rest/RestAuthenticationDetails.java

@ -29,18 +29,10 @@ public class RestAuthenticationDetails implements Serializable {
private final Client userAgent;
public RestAuthenticationDetails(HttpServletRequest request) {
this.clientAddress = getClientIP(request);
this.clientAddress = request.getRemoteAddr();
this.userAgent = getUserAgent(request);
}
private static String getClientIP(HttpServletRequest request) {
String xfHeader = request.getHeader("X-Forwarded-For");
if (xfHeader == null) {
return request.getRemoteAddr();
}
return xfHeader.split(",")[0];
}
private static Client getUserAgent(HttpServletRequest request) {
Parser uaParser = new Parser();
return uaParser.parse(request.getHeader("User-Agent"));

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

@ -197,6 +197,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Statistics reporting interval, set to send summarized data every 10 seconds by default
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
check:
@ -1910,6 +1912,9 @@ queue:
print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}"
# Max length of the error message that is printed by statistics
max-error-message-length: "${TB_QUEUE_RULE_ENGINE_MAX_ERROR_MESSAGE_LENGTH:4096}"
prometheus-stats:
# Enable/disable Prometheus statistics for individual Rule Engine message processing (records time in ms for success/failure).
enabled: "${TB_QUEUE_RULE_ENGINE_PROMETHEUS_STATS_ENABLED:false}"
# After a queue is deleted (or the profile's isolation option was disabled), Rule Engine will continue reading related topics during this period before deleting the actual topics
topic-deletion-delay: "${TB_QUEUE_RULE_ENGINE_TOPIC_DELETION_DELAY_SEC:15}"
# Size of the thread pool that handles such operations as partition changes, config updates, queue deletion
@ -2036,7 +2041,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
health:
elasticsearch:
# Enable the org.springframework.boot.actuate.elasticsearch.ElasticsearchRestClientHealthIndicator.doHealthCheck

1
application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java

@ -430,6 +430,7 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest {
tenantProfileService.deleteTenantProfiles(TenantId.SYS_TENANT_ID);
jdbcTemplate.execute("TRUNCATE TABLE notification");
jdbcTemplate.execute("TRUNCATE TABLE audit_log");
log.debug("Executed web test teardown");
}

17
application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java

@ -407,20 +407,15 @@ public class TenantProfileControllerTest extends AbstractControllerTest {
testBroadcastEntityStateChangeEventNeverTenantProfile();
}
private void awaitAuditLog(String awaitMessage, TenantProfileId tenantProfileId, ActionType expectedAction) throws Exception {
private void awaitAuditLog(String awaitMessage, TenantProfileId tenantProfileId, ActionType expectedAction) {
Awaitility.await(awaitMessage)
.atMost(TIMEOUT, TimeUnit.SECONDS)
.until(() ->
doGetTypedWithTimePageLink(
"/api/audit/logs/entity/TENANT_PROFILE/" + tenantProfileId.getId() + "?",
new TypeReference<PageData<AuditLog>>() {
},
new TimePageLink(5))
.getData()
.stream()
.anyMatch(log -> log.getActionType() == expectedAction)
.until(() -> doGetTypedWithTimePageLink(
"/api/audit/logs?",
new TypeReference<PageData<AuditLog>>() {},
new TimePageLink(100)).getData().stream()
.anyMatch(log -> log.getEntityId().equals(tenantProfileId) && log.getActionType() == expectedAction)
);
}
private TenantProfile createTenantProfile(String name) {

41
application/src/test/java/org/thingsboard/server/service/apiusage/ApiUsageTest.java

@ -19,6 +19,8 @@ import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.SaveDeviceWithCredentialsRequest;
@ -30,6 +32,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.stats.TbApiUsageReportClient;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.controller.TbUrlConstants;
import org.thingsboard.server.dao.service.DaoSqlTest;
@ -46,6 +49,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
@TestPropertySource(properties = {
"usage.stats.report.enabled=true",
"usage.stats.report.interval=2",
"usage.stats.report.urgent_interval=1",
"usage.stats.gauge_report_interval=1",
})
public class ApiUsageTest extends AbstractControllerTest {
@ -54,9 +58,12 @@ public class ApiUsageTest extends AbstractControllerTest {
private User tenantAdmin;
private static final int MAX_DP_ENABLE_VALUE = 12;
private static final int MAX_SMS_ENABLE_VALUE = 10;
private static final double WARN_THRESHOLD_VALUE = 0.5;
@Autowired
private ApiUsageStateService apiUsageStateService;
@Autowired
private TbApiUsageReportClient apiUsageReportClient;
@Before
public void beforeTest() throws Exception {
@ -82,7 +89,7 @@ public class ApiUsageTest extends AbstractControllerTest {
}
@Test
public void testTelemetryApiCall() throws Exception {
public void testDbStorageApiUsage() throws Exception {
Device device = createDevice();
assertNotNull(device);
String telemetryPayload = "{\"temperature\":25, \"humidity\":60}";
@ -94,7 +101,8 @@ public class ApiUsageTest extends AbstractControllerTest {
doPostAsync(url, telemetryPayload, String.class, status().isOk());
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> assertEquals(ApiUsageStateValue.WARNING, apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState()));
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(ApiUsageStateValue.WARNING, getUsageState().getDbStorageState()));
long VALUE_DISABLE = (long) (MAX_DP_ENABLE_VALUE - (MAX_DP_ENABLE_VALUE * WARN_THRESHOLD_VALUE)) / 2;
@ -104,10 +112,35 @@ public class ApiUsageTest extends AbstractControllerTest {
await().atMost(TIMEOUT, TimeUnit.SECONDS)
.untilAsserted(() -> {
assertEquals(ApiUsageStateValue.DISABLED, apiUsageStateService.findTenantApiUsageState(tenantId).getDbStorageState());
assertEquals(ApiUsageStateValue.DISABLED, getUsageState().getDbStorageState());
});
}
@Test
public void testSmsApiUsage() {
long smsWarnThreshold = (long) (MAX_SMS_ENABLE_VALUE * WARN_THRESHOLD_VALUE);
for (int i = 0; i < smsWarnThreshold; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.SMS_EXEC_COUNT);
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(ApiUsageStateValue.WARNING, getUsageState().getSmsExecState()));
long smsDisableCount = MAX_SMS_ENABLE_VALUE - smsWarnThreshold;
for (int i = 0; i < smsDisableCount; i++) {
apiUsageReportClient.report(tenantId, null, ApiUsageRecordKey.SMS_EXEC_COUNT);
}
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(ApiUsageStateValue.DISABLED, getUsageState().getSmsExecState()));
}
private ApiUsageState getUsageState() {
return apiUsageStateService.findTenantApiUsageState(tenantId);
}
private TenantProfile createTenantProfile() {
TenantProfile tenantProfile = new TenantProfile();
tenantProfile.setName("Tenant Profile");
@ -116,6 +149,8 @@ public class ApiUsageTest extends AbstractControllerTest {
TenantProfileData tenantProfileData = new TenantProfileData();
DefaultTenantProfileConfiguration config = DefaultTenantProfileConfiguration.builder()
.maxDPStorageDays(MAX_DP_ENABLE_VALUE)
.maxSms(MAX_SMS_ENABLE_VALUE)
.smsEnabled(true)
.warnThreshold(WARN_THRESHOLD_VALUE)
.build();

2
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/RelatedEntitiesAggregationCalculatedFieldStateTest.java

@ -234,6 +234,8 @@ public class RelatedEntitiesAggregationCalculatedFieldStateTest {
config.setUseLatestTs(true);
config.setScheduledUpdateInterval(10);
calculatedField.setConfiguration(config);
calculatedField.setVersion(1L);
return calculatedField;

260
application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java

@ -15,14 +15,17 @@
*/
package org.thingsboard.server.service.queue;
import lombok.extern.slf4j.Slf4j;
import org.junit.After;
import org.junit.Assert;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.junit.MockitoJUnitRunner;
import com.google.common.util.concurrent.MoreExecutors;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;
@ -35,30 +38,241 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.junit.Assert.assertTrue;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@Slf4j
@RunWith(MockitoJUnitRunner.class)
public class TbMsgPackProcessingContextTest {
@ExtendWith(MockitoExtension.class)
class TbMsgPackProcessingContextTest {
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
@Mock
TbRuleEngineSubmitStrategy submitStrategy;
@Mock
TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg> mockMsg;
ConcurrentMap<UUID, TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg>> pendingMap;
public static final int TIMEOUT = 10;
ExecutorService executorService;
@After
public void tearDown() {
@BeforeEach
void setup() {
pendingMap = new ConcurrentHashMap<>();
lenient().when(submitStrategy.getPendingMap()).thenReturn(pendingMap);
}
@AfterEach
void tearDown() {
if (executorService != null) {
executorService.shutdownNow();
MoreExecutors.shutdownAndAwaitTermination(executorService, 5, TimeUnit.SECONDS);
}
}
@Test
public void testHighConcurrencyCase() throws InterruptedException {
//log.warn("preparing the test...");
void testAwait_shouldReturnTrue_whenOnSuccessIsCalledBeforeTimeout() throws InterruptedException {
// GIVEN - a context with one pending message
executorService = Executors.newSingleThreadExecutor();
UUID msgId = UUID.randomUUID();
pendingMap.put(msgId, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
// WHEN - onSuccess() is called in another thread before timeout
executorService.submit(() -> {
try {
Thread.sleep(100);
context.onSuccess(msgId);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// THEN - await() should return true (successful completion)
boolean result = context.await(5000, TimeUnit.MILLISECONDS);
assertThat(result).as("await() should return true when latch is counted down before timeout").isTrue();
// Verify the message was moved to success map
assertThat(context.getSuccessMap()).containsKey(msgId);
assertThat(context.getPendingMap()).isEmpty();
assertThat(context.getExceptionsMap()).isEmpty();
// Verify submit strategy was notified about successful message processing
then(submitStrategy).should().onSuccess(msgId);
}
@Test
void testAwait_shouldReturnTrue_whenOnFailureIsCalledBeforeTimeout() throws InterruptedException {
// GIVEN - a context with one pending message
executorService = Executors.newSingleThreadExecutor();
UUID msgId = UUID.randomUUID();
pendingMap.put(msgId, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
var exception = new RuleEngineException("Test exception");
// WHEN - onFailure() is called in another thread before timeout
executorService.submit(() -> {
try {
Thread.sleep(100);
context.onFailure(tenantId, msgId, exception);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// THEN - await() should return true (successful completion, even if message processing failed)
boolean result = context.await(5000, TimeUnit.MILLISECONDS);
assertThat(result).as("await() should return true when latch is counted down before timeout").isTrue();
// Verify the exception was added to exceptions map
assertThat(context.getSuccessMap()).isEmpty();
assertThat(context.getPendingMap()).isEmpty();
assertThat(context.getExceptionsMap()).containsEntry(tenantId, exception);
}
@Test
void testAwait_shouldReturnFalse_whenTimeoutOccurs() throws InterruptedException {
// GIVEN - a context with one pending message and no processing
UUID msgId = UUID.randomUUID();
pendingMap.put(msgId, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
// WHEN - await() is called with short timeout and no message processing happens
long startTime = System.nanoTime();
boolean result = context.await(100, TimeUnit.MILLISECONDS);
long elapsedTime = System.nanoTime() - startTime;
// THEN - await() should return false (timeout occurred)
assertThat(result).as("await() should return false when timeout occurs").isFalse();
assertThat(elapsedTime).as("await() should wait for at least the timeout duration").isGreaterThanOrEqualTo(100L);
// Message should still be in pending map
assertThat(context.getSuccessMap()).isEmpty();
assertThat(context.getPendingMap()).containsKey(msgId);
assertThat(context.getExceptionsMap()).isEmpty();
}
@Test
void testAwait_shouldHandleMultiplePendingMessages() throws InterruptedException {
// GIVEN - a context with multiple pending messages
executorService = Executors.newSingleThreadExecutor();
UUID msgId1 = UUID.randomUUID();
UUID msgId2 = UUID.randomUUID();
UUID msgId3 = UUID.randomUUID();
pendingMap.put(msgId1, mockMsg);
pendingMap.put(msgId2, mockMsg);
pendingMap.put(msgId3, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
// WHEN - messages are processed one by one
executorService.submit(() -> {
try {
Thread.sleep(50);
context.onSuccess(msgId1);
Thread.sleep(50);
context.onSuccess(msgId2);
Thread.sleep(50);
context.onSuccess(msgId3);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// THEN - await() should return true only after all messages are processed
boolean result = context.await(5000, TimeUnit.MILLISECONDS);
assertThat(result).as("await() should return true after all messages are processed").isTrue();
// All messages should be in success map
assertThat(context.getSuccessMap()).containsKeys(msgId1, msgId2, msgId3);
assertThat(context.getPendingMap()).isEmpty();
assertThat(context.getExceptionsMap()).isEmpty();
}
@Test
void testAwait_shouldNotCountDownPrematurely_withMultipleMessages() throws InterruptedException {
// GIVEN - a context with multiple pending messages
executorService = Executors.newSingleThreadExecutor();
UUID msgId1 = UUID.randomUUID();
UUID msgId2 = UUID.randomUUID();
pendingMap.put(msgId1, mockMsg);
pendingMap.put(msgId2, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
// WHEN - only one message is processed
executorService.submit(() -> {
try {
Thread.sleep(100);
context.onSuccess(msgId1);
// msgId2 still in processing
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// THEN: await should timeout because not all messages were processed
boolean result = context.await(2000, TimeUnit.MILLISECONDS);
assertThat(result).as("await() should timeout when not all messages are processed").isFalse();
// One message in success, one still pending
assertThat(context.getSuccessMap()).containsOnlyKeys(msgId1);
assertThat(context.getPendingMap()).containsOnlyKeys(msgId2);
assertThat(context.getExceptionsMap()).isEmpty();
}
@Test
void testAwait_shouldHandleMixedSuccessAndFailure() throws InterruptedException {
// GIVEN - multiple messages
executorService = Executors.newSingleThreadExecutor();
UUID msgId1 = UUID.randomUUID();
UUID msgId2 = UUID.randomUUID();
pendingMap.put(msgId1, mockMsg);
pendingMap.put(msgId2, mockMsg);
var context = new TbMsgPackProcessingContext("test-queue", submitStrategy, false);
var exception = new RuleEngineException("Test exception");
// WHEN - one succeeds, one fails
executorService.submit(() -> {
try {
Thread.sleep(50);
context.onSuccess(msgId1);
Thread.sleep(50);
context.onFailure(tenantId, msgId2, exception);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
// THEN - await() should complete successfully
boolean result = context.await(5000, TimeUnit.MILLISECONDS);
assertThat(result).as("await() should return true when all messages are processed").isTrue();
assertThat(context.getSuccessMap()).containsOnlyKeys(msgId1);
assertThat(context.getPendingMap()).isEmpty();
assertThat(context.getExceptionsMap()).containsEntry(tenantId, exception);
}
@Test
void testHighConcurrencyCase() throws InterruptedException {
int msgCount = 1000;
int parallelCount = 5;
executorService = Executors.newFixedThreadPool(parallelCount, ThingsBoardThreadFactory.forName(getClass().getSimpleName() + "-test-scope"));
@ -76,28 +290,24 @@ public class TbMsgPackProcessingContextTest {
final CountDownLatch startLatch = new CountDownLatch(1);
final CountDownLatch finishLatch = new CountDownLatch(parallelCount);
for (int i = 0; i < parallelCount; i++) {
//final String taskName = "" + uuid + " " + i;
executorService.submit(() -> {
//log.warn("ready {}", taskName);
readyLatch.countDown();
try {
startLatch.await();
} catch (InterruptedException e) {
Assert.fail("failed to await");
fail("failed to await");
}
//log.warn("go {}", taskName);
context.onSuccess(uuid);
finishLatch.countDown();
});
}
assertTrue(readyLatch.await(TIMEOUT, TimeUnit.SECONDS));
assertTrue(readyLatch.await(10, TimeUnit.SECONDS));
Thread.yield();
startLatch.countDown(); //run all-at-once submitted tasks
assertTrue(finishLatch.await(TIMEOUT, TimeUnit.SECONDS));
assertTrue(finishLatch.await(10, TimeUnit.SECONDS));
}
assertTrue(context.await(TIMEOUT, TimeUnit.SECONDS));
assertTrue(context.await(10, TimeUnit.SECONDS));
verify(strategyMock, times(msgCount)).onSuccess(any(UUID.class));
}
}

281
application/src/test/java/org/thingsboard/server/service/queue/ruleengine/RuleEngineConsumerLoopTest.java

@ -0,0 +1,281 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.queue.ruleengine;
import com.google.common.util.concurrent.MoreExecutors;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InOrder;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.queue.ProcessingStrategy;
import org.thingsboard.server.common.data.queue.ProcessingStrategyType;
import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.data.queue.SubmitStrategy;
import org.thingsboard.server.common.data.queue.SubmitStrategyType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.TbQueueAdmin;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.common.DefaultTbQueueMsgHeaders;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.memory.DefaultInMemoryStorage;
import org.thingsboard.server.queue.memory.InMemoryStorage;
import org.thingsboard.server.queue.memory.InMemoryTbQueueConsumer;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContext;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContextFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory;
import org.thingsboard.server.service.stats.RuleEngineStatisticsService;
import java.time.Duration;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class RuleEngineConsumerLoopTest {
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
DeviceId deviceId = new DeviceId(UUID.randomUUID());
InMemoryStorage storage;
@Mock
ActorSystemContext actorContext;
@Mock
StatsFactory statsFactory;
@Mock
TbRuleEngineQueueFactory queueFactory;
@Mock
RuleEngineStatisticsService statisticsService;
@Mock
TbServiceInfoProvider serviceInfoProvider;
@Mock
PartitionService partitionService;
@Mock
TbQueueProducerProvider producerProvider;
@Mock
TbQueueAdmin queueAdmin;
@Mock
TbMsgPackProcessingContextFactory packProcessingContextFactory;
@Mock
TbMsgPackProcessingContext packCtx;
Queue mainQueue;
TbQueueConsumer<TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg>> consumer;
TbRuleEngineConsumerContext ruleEngineConsumerContext;
TbRuleEngineQueueConsumerManager consumerManager;
ExecutorService consumersExecutor;
ScheduledExecutorService scheduler;
ExecutorService mgmtExecutor;
@BeforeEach
void setup() throws InterruptedException {
consumersExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("tb-rule-engine-consumer"));
scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("tb-rule-engine-consumer-scheduler");
mgmtExecutor = ThingsBoardExecutors.newWorkStealingPool(1, "tb-rule-engine-mgmt");
mainQueue = new Queue();
mainQueue.setTenantId(TenantId.SYS_TENANT_ID);
mainQueue.setName("Main");
mainQueue.setTopic("tb_rule_engine.main");
mainQueue.setPollInterval(25);
mainQueue.setPartitions(1);
mainQueue.setConsumerPerPartition(false);
mainQueue.setPackProcessingTimeout(2000L);
var submitStrategy = new SubmitStrategy();
submitStrategy.setType(SubmitStrategyType.BURST);
submitStrategy.setBatchSize(1000);
mainQueue.setSubmitStrategy(submitStrategy);
var processingStrategy = new ProcessingStrategy();
processingStrategy.setType(ProcessingStrategyType.SKIP_ALL_FAILURES);
processingStrategy.setRetries(3);
processingStrategy.setFailurePercentage(0.0);
processingStrategy.setPauseBetweenRetries(3);
processingStrategy.setMaxPauseBetweenRetries(3);
mainQueue.setProcessingStrategy(processingStrategy);
storage = new DefaultInMemoryStorage();
consumer = spy(new InMemoryTbQueueConsumer<>(storage, mainQueue.getTopic()));
given(queueFactory.createToRuleEngineMsgConsumer(eq(mainQueue), isNull())).willReturn(consumer);
ruleEngineConsumerContext = new TbRuleEngineConsumerContext(
actorContext, statsFactory, new TbRuleEngineSubmitStrategyFactory(), new TbRuleEngineProcessingStrategyFactory(),
queueFactory, statisticsService, serviceInfoProvider, partitionService, producerProvider, queueAdmin
);
ruleEngineConsumerContext.setPollDuration(25);
ruleEngineConsumerContext.setPackProcessingTimeout(2000);
ruleEngineConsumerContext.setStatsEnabled(false); // true by default
ruleEngineConsumerContext.setPrometheusStatsEnabled(false);
ruleEngineConsumerContext.setTopicDeletionDelayInSec(15);
ruleEngineConsumerContext.setMgmtThreadPoolSize(12);
// Tell the (mock) context factory to return (mock) message pack context
given(packProcessingContextFactory.create(
eq(mainQueue.getName()),
any(TbRuleEngineSubmitStrategy.class),
eq(false)
)).willAnswer(invocation -> {
TbRuleEngineSubmitStrategy realStrategy = invocation.getArgument(1);
when(packCtx.getPendingMap()).thenAnswer(i -> realStrategy.getPendingMap());
when(packCtx.getFailedMap()).thenReturn(new ConcurrentHashMap<>());
return packCtx;
});
// Tell the (mock) context's await() to return 'false' (always timeout) immediately
given(packCtx.await(anyLong(), any(TimeUnit.class))).willReturn(false);
consumerManager = TbRuleEngineQueueConsumerManager.create()
.ctx(ruleEngineConsumerContext)
.queueKey(new QueueKey(ServiceType.TB_RULE_ENGINE, mainQueue))
.consumerExecutor(consumersExecutor)
.scheduler(scheduler)
.taskExecutor(mgmtExecutor)
.packProcessingContextFactory(packProcessingContextFactory)
.build();
}
@AfterEach
void destroy() {
MoreExecutors.shutdownAndAwaitTermination(scheduler, Duration.ofSeconds(30));
MoreExecutors.shutdownAndAwaitTermination(mgmtExecutor, Duration.ofSeconds(30));
MoreExecutors.shutdownAndAwaitTermination(consumersExecutor, Duration.ofSeconds(30));
}
@Test
void consumerLoopTest_verifyOperationsOrder() throws InterruptedException {
// Create partition
var partition = TopicPartitionInfo.builder()
.tenantId(TenantId.SYS_TENANT_ID)
.topic(mainQueue.getTopic())
.partition(0)
.myPartition(true)
.useInternalPartition(false)
.build();
// Put 10k messages to the queue
for (int i = 0; i < 10_000; i++) {
var tbMsg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(deviceId)
.data("{\"temperature\":123}")
.metaData(TbMsgMetaData.EMPTY)
.build();
var toRuleEngineMsg = TransportProtos.ToRuleEngineMsg.newBuilder()
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.setTenantIdMSB(tenantId.getId().getMostSignificantBits())
.setTbMsgProto(TbMsg.toProto(tbMsg))
.addAllRelationTypes(Set.of("Success"))
.build();
storage.put(partition.getFullTopicName(), new TbProtoQueueMsg<>(UUID.randomUUID(), toRuleEngineMsg, new DefaultTbQueueMsgHeaders()));
}
// Count how many polls were made
var totalPolls = new AtomicInteger(0);
var emptyPolls = new AtomicInteger(0);
doAnswer(invocation -> {
totalPolls.incrementAndGet();
@SuppressWarnings("unchecked")
var messages = (List<TbQueueMsg>) invocation.callRealMethod();
if (messages.isEmpty()) {
emptyPolls.incrementAndGet();
}
return messages;
}).when(consumer).poll(mainQueue.getPollInterval());
// Count how many commits were made
var totalCommits = new AtomicInteger(0);
doAnswer(invocation -> {
totalCommits.incrementAndGet();
return invocation.callRealMethod();
}).when(consumer).commit();
// Initialize consumer
consumerManager.init(mainQueue);
// Assign partition to the consumer
consumerManager.update(Set.of(partition));
// Give some time for the consumer to get all messages
await().atMost(Duration.ofSeconds(10L)).until(() -> storage.getLagTotal() == 0);
// Stop consumer
consumerManager.stop();
consumerManager.awaitStop();
// Determine number of non-empty consumer iterations made, since polling does not stop immediately after consuming all messages and may do a few empty polls
int nonEmptyPolls = totalPolls.get() - emptyPolls.get();
// Verify that there is 10 polls and 10 matching commits
// Each poll consumes 1k messages and queue has 10k total, so that means 10k total msgs / 1k msgs per poll = 10 polls
assertThat(nonEmptyPolls).isEqualTo(10).isEqualTo(totalCommits.get());
// Verify that poll-await-commit cycle happened in order with correct await timeout
InOrder inOrder = inOrder(consumer, packCtx);
for (int i = 0; i < nonEmptyPolls; i++) {
inOrder.verify(consumer).poll(mainQueue.getPollInterval());
inOrder.verify(packCtx).await(mainQueue.getPackProcessingTimeout(), TimeUnit.MILLISECONDS);
inOrder.verify(consumer).commit();
}
}
}

2
application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java

@ -59,6 +59,7 @@ import org.thingsboard.server.queue.provider.KafkaMonolithQueueFactory;
import org.thingsboard.server.queue.provider.KafkaTbRuleEngineQueueFactory;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContextFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory;
import org.thingsboard.server.service.stats.RuleEngineStatisticsService;
@ -194,6 +195,7 @@ public class TbRuleEngineQueueConsumerManagerTest {
.consumerExecutor(consumersExecutor)
.scheduler(scheduler)
.taskExecutor(mgmtExecutor)
.packProcessingContextFactory(new TbMsgPackProcessingContextFactory.DefaultTbMsgPackProcessingContextFactory())
.build();
}

2
application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineStrategyTest.java

@ -45,6 +45,7 @@ import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerTask.ConsumerKey;
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.service.queue.TbMsgPackProcessingContextFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory;
@ -196,6 +197,7 @@ public class TbRuleEngineStrategyTest {
var consumerManager = TbRuleEngineQueueConsumerManager.create()
.ctx(ruleEngineConsumerContext)
.queueKey(queueKey)
.packProcessingContextFactory(new TbMsgPackProcessingContextFactory.DefaultTbMsgPackProcessingContextFactory())
.build();
consumerManager.init(queue);

1
application/src/test/java/org/thingsboard/server/service/sms/DefaultSmsServiceTest.java

@ -52,6 +52,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
@TestPropertySource(properties = {
"usage.stats.report.enabled=true",
"usage.stats.report.interval=1",
"usage.stats.report.urgent_interval=1"
})
public class DefaultSmsServiceTest extends AbstractControllerTest {
@SpyBean

1
application/src/test/java/org/thingsboard/server/service/stats/DevicesStatisticsTest.java

@ -41,6 +41,7 @@ import static org.awaitility.Awaitility.await;
@TestPropertySource(properties = {
"usage.stats.report.enabled=true",
"usage.stats.report.interval=2",
"usage.stats.report.urgent_interval=1",
"usage.stats.gauge_report_interval=1",
"usage.stats.devices.report_interval=3",
"state.defaultStateCheckIntervalInSec=3",

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

@ -140,45 +140,61 @@ public class FwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
}
private void startDownloading() {
long delay = 0;
// Step 1: state = 1
scheduler.schedule(() -> {
try {
state.set(1);
fireResourceChange(3);
Thread.sleep(100);
state.set(2);
fireResourceChange(3);
} catch (Exception e) {
}
}, 100, TimeUnit.MILLISECONDS);
state.set(1);
fireResourceChange(3);
log.info("Downloading started: state=[{}]", state.get());
}, delay, TimeUnit.MILLISECONDS);
delay += 100; // next step after 100 ms
// Step 2: state = 2
scheduler.schedule(() -> {
state.set(2);
fireResourceChange(3);
log.info("Downloading in progress: state=[{}]", state.get());
}, delay, TimeUnit.MILLISECONDS);
}
private void startUpdating(LwM2mServer identity) {
scheduler.schedule(() -> {
try {
// Update state + result
state.set(3);
fireResourceChange(3);
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;
// Delayed reset pkgName/pkgVersion, after reboot + registration
scheduler.schedule(() -> {
this.pkgName = this.pkgNameDef;
fireResourceChange(6);
this.pkgVersion = this.pkgVersionDef;
fireResourceChange(7);
log.info("FW resources updating to new values: pkgName=[{}], pkgVersion=[{}]",
this.pkgName, this.pkgVersion);
}, 15, TimeUnit.SECONDS); // 15 sec — safe timing
}
} catch (Exception e) {
log.error("Error during firmware update", e);
}
}, 100, TimeUnit.MILLISECONDS);
}, 0, TimeUnit.SECONDS); // start immediately, without further delay
}
protected void setLeshanClient(LeshanClient leshanClient) {
this.leshanClient = leshanClient;
}
}

60
application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SwLwM2MDevice.java

@ -33,6 +33,8 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.thingsboard.server.controller.AbstractWebTest.TIMEOUT;
@Slf4j
public class SwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
@ -85,10 +87,7 @@ public class SwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
log.info("Write on Device resource /{}/{}/{}", getModel().id, getId(), resourceId);
switch (resourceId) {
case 2:
startDownloading();
return WriteResponse.success();
case 3:
case 2, 3:
startDownloading();
return WriteResponse.success();
default:
@ -123,25 +122,34 @@ public class SwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
}
private void startDownloading() {
long delay = 0;
// Step 1: start downloading
scheduler.schedule(() -> {
try {
state.set(1);
updateResult.set(1);
fireResourceChange(7);
fireResourceChange(9);
Thread.sleep(100);
state.set(2);
fireResourceChange(7);
Thread.sleep(100);
state.set(3);
fireResourceChange(7);
Thread.sleep(100);
updateResult.set(3);
fireResourceChange(9);
} catch (Exception e) {
}
}, 100, TimeUnit.MILLISECONDS);
state.set(1);
updateResult.set(1);
fireResourceChange(7);
fireResourceChange(9);
}, delay, TimeUnit.MILLISECONDS);
delay += 100;
// Step 2: downloading in progress
scheduler.schedule(() -> {
state.set(2);
fireResourceChange(7);
}, delay, TimeUnit.MILLISECONDS);
delay += 100;
// Step 3: downloading finished
scheduler.schedule(() -> {
state.set(3);
fireResourceChange(7);
updateResult.set(3);
fireResourceChange(9);
}, delay, TimeUnit.MILLISECONDS);
}
private void startUpdating() {
@ -150,7 +158,13 @@ public class SwLwM2MDevice extends BaseInstanceEnabler implements Destroyable {
updateResult.set(2);
fireResourceChange(7);
fireResourceChange(9);
// Optional: delayed log about FW update
scheduler.schedule(() -> {
log.info("FW resources updating to new values: state=[{}], updateResult=[{}]",
state.get(), updateResult.get());
}, 500, TimeUnit.MILLISECONDS);
}, 100, TimeUnit.MILLISECONDS);
}
}

13
common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageRecordKey.java

@ -25,8 +25,8 @@ public enum ApiUsageRecordKey {
RE_EXEC_COUNT(ApiFeature.RE, "ruleEngineExecutionCount", "ruleEngineExecutionLimit", "Rule Engine execution"),
JS_EXEC_COUNT(ApiFeature.JS, "jsExecutionCount", "jsExecutionLimit", "JavaScript execution"),
TBEL_EXEC_COUNT(ApiFeature.TBEL, "tbelExecutionCount", "tbelExecutionLimit", "Tbel execution"),
EMAIL_EXEC_COUNT(ApiFeature.EMAIL, "emailCount", "emailLimit", "email message"),
SMS_EXEC_COUNT(ApiFeature.SMS, "smsCount", "smsLimit", "SMS message"),
EMAIL_EXEC_COUNT(ApiFeature.EMAIL, "emailCount", "emailLimit", "email message", true, true),
SMS_EXEC_COUNT(ApiFeature.SMS, "smsCount", "smsLimit", "SMS message", true, true),
CREATED_ALARMS_COUNT(ApiFeature.ALARM, "createdAlarmsCount", "createdAlarmsLimit", "alarm"),
ACTIVE_DEVICES("activeDevicesCount"),
INACTIVE_DEVICES("inactiveDevicesCount");
@ -50,21 +50,24 @@ public enum ApiUsageRecordKey {
private final String unitLabel;
@Getter
private final boolean counter;
@Getter
private final boolean urgent; // urgent keys are reported at a shorter interval for quicker usage state updates
ApiUsageRecordKey(ApiFeature apiFeature, String apiCountKey, String apiLimitKey, String unitLabel) {
this(apiFeature, apiCountKey, apiLimitKey, unitLabel, true);
this(apiFeature, apiCountKey, apiLimitKey, unitLabel, true, false);
}
ApiUsageRecordKey(String apiCountKey) {
this(null, apiCountKey, null, null, false);
this(null, apiCountKey, null, null, false, false);
}
ApiUsageRecordKey(ApiFeature apiFeature, String apiCountKey, String apiLimitKey, String unitLabel, boolean counter) {
ApiUsageRecordKey(ApiFeature apiFeature, String apiCountKey, String apiLimitKey, String unitLabel, boolean counter, boolean urgent) {
this.apiFeature = apiFeature;
this.apiCountKey = apiCountKey;
this.apiLimitKey = apiLimitKey;
this.unitLabel = unitLabel;
this.counter = counter;
this.urgent = urgent;
}
public static ApiUsageRecordKey[] getKeys(ApiFeature feature) {

8
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java

@ -22,14 +22,14 @@ public interface ScheduledUpdateSupportedCalculatedFieldConfiguration extends Ca
boolean isScheduledUpdateEnabled();
@PositiveOrZero
int getScheduledUpdateInterval();
Integer getScheduledUpdateInterval();
void setScheduledUpdateInterval(int interval);
void setScheduledUpdateInterval(Integer interval);
default void validate(long minAllowedScheduledUpdateInterval) {
if (getScheduledUpdateInterval() < minAllowedScheduledUpdateInterval) {
throw new IllegalArgumentException("Scheduled update interval is less than configured " +
"minimum allowed interval in tenant profile: " + minAllowedScheduledUpdateInterval);
throw new IllegalArgumentException("Scheduled update interval (" + getScheduledUpdateInterval() +
" seconds) is less than minimum allowed interval in tenant profile: " + minAllowedScheduledUpdateInterval + " seconds");
}
}
}

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

@ -45,7 +45,7 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A
private Output output;
private boolean useLatestTs;
private int scheduledUpdateInterval;
private Integer scheduledUpdateInterval;
@Override
public CalculatedFieldType getType() {

5
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java

@ -48,7 +48,7 @@ public class GeofencingCalculatedFieldConfiguration implements ArgumentsBasedCal
private Map<String, ZoneGroupConfiguration> zoneGroups;
private boolean scheduledUpdateEnabled;
private int scheduledUpdateInterval;
private Integer scheduledUpdateInterval;
@NotNull
private Output output;
@ -88,6 +88,9 @@ public class GeofencingCalculatedFieldConfiguration implements ArgumentsBasedCal
@Override
public void validate() {
if (scheduledUpdateEnabled && scheduledUpdateInterval == null) {
throw new IllegalArgumentException("Refresh interval is required when periodic zone group refresh is enabled.");
}
zoneGroups.forEach((key, value) -> value.validate(key));
}

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

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.tenant.profile;
import io.swagger.v3.oas.annotations.media.Schema;
import jakarta.validation.constraints.Positive;
import jakarta.validation.constraints.PositiveOrZero;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
@ -173,6 +174,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
@Schema(example = "10")
private long maxArgumentsPerCF = 10;
@Schema(example = "10")
@PositiveOrZero
private int minAllowedScheduledUpdateIntervalInSecForCF = 10;
@Builder.Default
@Schema(example = "2")

4
common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfigurationTest.java

@ -47,8 +47,8 @@ public class ScheduledUpdateSupportedCalculatedFieldConfigurationTest {
assertThatThrownBy(() -> cfg.validate(minAllowedInterval))
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Scheduled update interval is less than configured " +
"minimum allowed interval in tenant profile: " + minAllowedInterval);
.hasMessage("Scheduled update interval (1 seconds) is less than " +
"minimum allowed interval in tenant profile: " + minAllowedInterval + " seconds");
}
}

14
common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java

@ -28,6 +28,7 @@ import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY;
@ -103,4 +104,17 @@ public class GeofencingCalculatedFieldConfigurationTest {
assertThat(allowedZonesArgument.getRefEntityKey()).isEqualTo(new ReferencedEntityKey("perimeter", ArgumentType.ATTRIBUTE, AttributeScope.SERVER_SCOPE));
}
@Test
void validateShouldThrowWhenScheduledUpdateEnabledButIntervalNotSet() {
var cfg = new GeofencingCalculatedFieldConfiguration();
cfg.setEntityCoordinates(mock(EntityCoordinates.class));
cfg.setZoneGroups(Map.of("zone", mock(ZoneGroupConfiguration.class)));
cfg.setScheduledUpdateEnabled(true);
cfg.setScheduledUpdateInterval(null);
assertThatThrownBy(cfg::validate)
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Refresh interval is required when periodic zone group refresh is enabled.");
}
}

8
common/edge-api/src/main/proto/edge.proto

@ -46,12 +46,12 @@ enum EdgeVersion {
V_4_2_0 = 12;
V_4_3_0 = 13;
V_4_2_1_2 = 14;
V_4_2_1_3 = 420;
V_4_2_2 = 4220;
V_4_3_0_1 = 15;
V_4_3_0_2 = 430;
V_4_4_0 = 440;
V_4_3_1 = 4310;
V_4_4_0 = 4400;
V_LATEST = 999;
V_LATEST = 99999;
}
/**

4
common/proto/src/main/java/org/thingsboard/server/common/adaptor/ProtoConverter.java

@ -156,11 +156,7 @@ public class ProtoConverter {
case BOOLEAN_V:
case LONG_V:
case DOUBLE_V:
break;
case STRING_V:
if (StringUtils.isEmpty(keyValueProto.getStringV())) {
throw new IllegalArgumentException("Value is empty for key: " + key + "!");
}
break;
case JSON_V:
try {

20
common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java

@ -63,8 +63,10 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
private boolean enabled;
@Value("${usage.stats.report.enabled_per_customer:false}")
private boolean enabledPerCustomer;
@Value("${usage.stats.report.interval:10}")
@Value("${usage.stats.report.interval:60}")
private int interval;
@Value("${usage.stats.report.urgent_interval:10}")
private int urgentInterval;
@Value("${usage.stats.report.pack_size:1024}")
private int packSize;
@ -83,20 +85,30 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient {
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) {
stats.put(key, new ConcurrentHashMap<>());
}
Random random = new Random();
scheduler.scheduleWithFixedDelay(() -> {
try {
reportStats();
reportStats(false);
} catch (Exception e) {
log.warn("Failed to report statistics: ", e);
}
}, new Random().nextInt(interval), interval, TimeUnit.SECONDS);
}, random.nextInt(interval), interval, TimeUnit.SECONDS);
scheduler.scheduleWithFixedDelay(() -> {
try {
reportStats(true);
} catch (Exception e) {
log.warn("Failed to report urgent statistics: ", e);
}
}, random.nextInt(urgentInterval), urgentInterval, TimeUnit.SECONDS);
}
}
private void reportStats() {
private void reportStats(boolean urgent) {
ConcurrentMap<ParentEntity, UsageStatsServiceMsg.Builder> report = new ConcurrentHashMap<>();
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) {
if (key.isUrgent() != urgent) continue;
ConcurrentMap<ReportLevel, AtomicLong> statsForKey = stats.get(key);
statsForKey.forEach((reportLevel, statsValue) -> {
long value = statsValue.get();

14
dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

@ -26,18 +26,21 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldFilter;
import org.thingsboard.server.common.data.cf.CalculatedFieldInfo;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration;
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.HasId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent;
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.service.validator.CalculatedFieldDataValidator;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import java.util.EnumSet;
import java.util.List;
@ -62,6 +65,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
private final EntityService entityService;
private final CalculatedFieldDao calculatedFieldDao;
private final CalculatedFieldDataValidator calculatedFieldDataValidator;
private final ApiLimitService apiLimitService;
@Override
public CalculatedField save(CalculatedField calculatedField) {
@ -70,6 +74,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
@Override
public CalculatedField save(CalculatedField calculatedField, boolean doValidate) {
setConfigurationDefaults(calculatedField);
CalculatedField oldCalculatedField = null;
if (doValidate) {
oldCalculatedField = calculatedFieldDataValidator.validate(calculatedField, CalculatedField::getTenantId);
@ -79,6 +84,15 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
return doSave(calculatedField, oldCalculatedField);
}
private void setConfigurationDefaults(CalculatedField calculatedField) {
if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config
&& config.getScheduledUpdateInterval() == null) {
int minScheduledUpdateInterval = (int) apiLimitService.getLimit(
calculatedField.getTenantId(), DefaultTenantProfileConfiguration::getMinAllowedScheduledUpdateIntervalInSecForCF
);
config.setScheduledUpdateInterval(minScheduledUpdateInterval);
}
}
private CalculatedField doSave(CalculatedField calculatedField, CalculatedField oldCalculatedField) {
try {

4
dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java

@ -93,11 +93,13 @@ public abstract class AbstractServiceTest {
@Autowired
protected EntityServiceRegistry entityServiceRegistry;
protected Tenant tenant;
protected TenantId tenantId;
@Before
public void beforeAbstractService() {
tenantId = createTenant().getId();
tenant = createTenant();
tenantId = tenant.getId();
}
@After

232
dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java

@ -15,10 +15,12 @@
*/
package org.thingsboard.server.dao.service;
import org.apache.commons.lang3.RandomUtils;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
@ -29,6 +31,10 @@ import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey;
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration;
import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.TimeSeriesOutput;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput;
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.cf.configuration.geofencing.EntityCoordinates;
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.geofencing.ZoneGroupConfiguration;
@ -39,6 +45,7 @@ import org.thingsboard.server.common.data.relation.RelationPathLevel;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.exception.DataValidationException;
import java.util.ArrayList;
@ -59,6 +66,8 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
private DeviceService deviceService;
@Autowired
private TbTenantProfileCache tbTenantProfileCache;
@Autowired
private TenantProfileService tenantProfileService;
@Test
public void testSaveCalculatedField() {
@ -82,8 +91,6 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
assertThat(updatedCalculatedField.getName()).isEqualTo(savedCalculatedField.getName());
assertThat(updatedCalculatedField.getVersion()).isEqualTo(savedCalculatedField.getVersion() + 1);
calculatedFieldService.deleteCalculatedField(tenantId, savedCalculatedField.getId());
}
@Test
@ -113,11 +120,11 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
int min = tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedScheduledUpdateIntervalInSecForCF();
int valueFromConfig = min - 10;
// Enable scheduling with an interval below tenant min
cfg.setScheduledUpdateEnabled(true);
cfg.setScheduledUpdateInterval(valueFromConfig);
int invalidInterval = RandomUtils.insecure().randomInt(1, min);
cfg.setScheduledUpdateInterval(invalidInterval);
// Create & save Calculated Field
CalculatedField cf = new CalculatedField();
@ -131,8 +138,8 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
assertThatThrownBy(() -> calculatedFieldService.save(cf))
.isInstanceOf(DataValidationException.class)
.hasCauseInstanceOf(IllegalArgumentException.class)
.hasMessageStartingWith("Scheduled update interval is less than configured " +
"minimum allowed interval in tenant profile: ");
.hasMessage("Scheduled update interval (" + invalidInterval +
" seconds) is less than minimum allowed interval in tenant profile: " + min + " seconds");
}
@Test
@ -233,8 +240,67 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
int savedInterval = geofencingCalculatedFieldConfiguration.getScheduledUpdateInterval();
assertThat(savedInterval).isEqualTo(valueFromConfig);
}
calculatedFieldService.deleteCalculatedField(tenantId, saved.getId());
@Test
public void testSaveGeofencingCalculatedField_shouldAcceptZeroScheduledUpdateIntervalWhenTenantProfileAllows() {
// GIVEN
var device = createTestDevice();
// Store original value and update tenant profile to allow 0 as min scheduled update interval
TenantProfile tenantProfile = tenantProfileService.findTenantProfileById(tenantId, tenant.getTenantProfileId());
int originalMinScheduledUpdateInterval = tenantProfile.getDefaultProfileConfiguration().getMinAllowedScheduledUpdateIntervalInSecForCF();
tenantProfile.getDefaultProfileConfiguration().setMinAllowedScheduledUpdateIntervalInSecForCF(0);
tenantProfileService.saveTenantProfile(tenantId, tenantProfile);
tbTenantProfileCache.evict(tenantProfile.getId());
try {
// Build a valid Geofencing configuration
var cfg = new GeofencingCalculatedFieldConfiguration();
// Coordinates: TS_LATEST, no dynamic source
var entityCoordinates = new EntityCoordinates("latitude", "longitude");
cfg.setEntityCoordinates(entityCoordinates);
// Zone-group argument (ATTRIBUTE) — make it DYNAMIC so scheduling is enabled
var zoneGroupConfiguration = new ZoneGroupConfiguration("allowed", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false);
var dynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration();
dynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE)));
zoneGroupConfiguration.setRefDynamicSourceConfiguration(dynamicSourceConfiguration);
cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration));
// Enable scheduling with interval = 0
cfg.setScheduledUpdateEnabled(true);
cfg.setScheduledUpdateInterval(0);
// Create Calculated Field
var cf = new CalculatedField();
cf.setTenantId(tenantId);
cf.setEntityId(device.getId());
cf.setType(CalculatedFieldType.GEOFENCING);
cf.setName("GF zero scheduled update interval test");
cf.setConfigurationVersion(0);
cf.setConfiguration(cfg);
var out = new AttributesOutput();
out.setScope(AttributeScope.SERVER_SCOPE);
cfg.setOutput(out);
// WHEN
CalculatedField saved = calculatedFieldService.save(cf);
// THEN
assertThat(saved).isNotNull();
assertThat(saved.getConfiguration()).isInstanceOf(GeofencingCalculatedFieldConfiguration.class);
var savedConfig = (GeofencingCalculatedFieldConfiguration) saved.getConfiguration();
assertThat(savedConfig.getScheduledUpdateInterval()).isEqualTo(0);
} finally {
// Restore original tenant profile value
tenantProfile.getProfileConfiguration().orElseThrow().setMinAllowedScheduledUpdateIntervalInSecForCF(originalMinScheduledUpdateInterval);
tenantProfileService.saveTenantProfile(tenantId, tenantProfile);
tbTenantProfileCache.evict(tenantProfile.getId());
}
}
@Test
@ -254,8 +320,6 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
CalculatedField fetchedCalculatedField = calculatedFieldService.findById(tenantId, savedCalculatedField.getId());
assertThat(fetchedCalculatedField).isEqualTo(savedCalculatedField);
calculatedFieldService.deleteCalculatedField(tenantId, savedCalculatedField.getId());
}
@Test
@ -267,6 +331,156 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
assertThat(calculatedFieldService.findById(tenantId, savedCalculatedField.getId())).isNull();
}
@Test
public void testSaveRelatedEntitiesAggregationCF_shouldUseMinScheduledUpdateIntervalFromTenantProfileWhenNotSet() {
// GIVEN
var device = createTestDevice();
var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration();
cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE));
var argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null));
cfg.setArguments(Map.of("temp", argument));
var metric = new AggMetric();
metric.setFunction(AggFunction.AVG);
metric.setInput(new AggKeyInput("temp"));
cfg.setMetrics(Map.of("avgTemp", metric));
var output = new TimeSeriesOutput();
output.setName("avgTemperature");
cfg.setOutput(output);
int minDeduplicationInterval = (int) tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedDeduplicationIntervalInSecForCF();
cfg.setDeduplicationIntervalInSec(minDeduplicationInterval);
// Do NOT set scheduledUpdateInterval - it should default to tenant profile min value
var cf = new CalculatedField();
cf.setTenantId(tenantId);
cf.setEntityId(device.getId());
cf.setType(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION);
cf.setName("Related Entities Aggregation CF - default scheduled interval test");
cf.setConfigurationVersion(0);
cf.setConfiguration(cfg);
// WHEN
CalculatedField saved = calculatedFieldService.save(cf);
// THEN
assertThat(saved).isNotNull();
assertThat(saved.getConfiguration()).isInstanceOf(RelatedEntitiesAggregationCalculatedFieldConfiguration.class);
var savedConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) saved.getConfiguration();
int expectedMinScheduledUpdateInterval = tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedScheduledUpdateIntervalInSecForCF();
assertThat(savedConfig.getScheduledUpdateInterval()).isEqualTo(expectedMinScheduledUpdateInterval);
}
@Test
public void testSaveRelatedEntitiesAggregationCF_shouldThrowWhenScheduledUpdateIntervalLessThanMinAllowed() {
// GIVEN
var device = createTestDevice();
var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration();
cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE));
var argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null));
cfg.setArguments(Map.of("temp", argument));
var metric = new AggMetric();
metric.setFunction(AggFunction.AVG);
metric.setInput(new AggKeyInput("temp"));
cfg.setMetrics(Map.of("avgTemp", metric));
var output = new TimeSeriesOutput();
output.setName("avgTemperature");
cfg.setOutput(output);
int minDeduplicationInterval = (int) tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedDeduplicationIntervalInSecForCF();
cfg.setDeduplicationIntervalInSec(minDeduplicationInterval);
int minScheduledUpdateInterval = tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedScheduledUpdateIntervalInSecForCF();
int invalidInterval = RandomUtils.insecure().randomInt(1, minScheduledUpdateInterval);
cfg.setScheduledUpdateInterval(invalidInterval);
var cf = new CalculatedField();
cf.setTenantId(tenantId);
cf.setEntityId(device.getId());
cf.setType(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION);
cf.setName("Related Entities Aggregation CF - invalid scheduled interval test");
cf.setConfigurationVersion(0);
cf.setConfiguration(cfg);
// WHEN-THEN
assertThatThrownBy(() -> calculatedFieldService.save(cf))
.isInstanceOf(DataValidationException.class)
.hasCauseInstanceOf(IllegalArgumentException.class)
.hasMessage("Scheduled update interval (" + invalidInterval +
" seconds) is less than minimum allowed interval in tenant profile: " + minScheduledUpdateInterval + " seconds");
}
@Test
public void testSaveRelatedEntitiesAggregationCF_shouldAcceptValidScheduledUpdateInterval() {
// GIVEN
var device = createTestDevice();
var cfg = new RelatedEntitiesAggregationCalculatedFieldConfiguration();
cfg.setRelation(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE));
var argument = new Argument();
argument.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null));
cfg.setArguments(Map.of("temp", argument));
var metric = new AggMetric();
metric.setFunction(AggFunction.AVG);
metric.setInput(new AggKeyInput("temp"));
cfg.setMetrics(Map.of("avgTemp", metric));
var output = new TimeSeriesOutput();
output.setName("avgTemperature");
cfg.setOutput(output);
int minDeduplicationInterval = (int) tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedDeduplicationIntervalInSecForCF();
cfg.setDeduplicationIntervalInSec(minDeduplicationInterval);
int minScheduledUpdateInterval = tbTenantProfileCache.get(tenantId)
.getDefaultProfileConfiguration()
.getMinAllowedScheduledUpdateIntervalInSecForCF();
int customScheduledUpdateInterval = minScheduledUpdateInterval + 100;
cfg.setScheduledUpdateInterval(customScheduledUpdateInterval);
var cf = new CalculatedField();
cf.setTenantId(tenantId);
cf.setEntityId(device.getId());
cf.setType(CalculatedFieldType.RELATED_ENTITIES_AGGREGATION);
cf.setName("Related Entities Aggregation CF - valid scheduled interval test");
cf.setConfigurationVersion(0);
cf.setConfiguration(cfg);
// WHEN
CalculatedField saved = calculatedFieldService.save(cf);
// THEN
assertThat(saved).isNotNull();
assertThat(saved.getConfiguration()).isInstanceOf(RelatedEntitiesAggregationCalculatedFieldConfiguration.class);
var savedConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) saved.getConfiguration();
assertThat(savedConfig.getScheduledUpdateInterval()).isEqualTo(customScheduledUpdateInterval);
}
private CalculatedField saveValidCalculatedField() {
Device device = createTestDevice();
CalculatedField calculatedField = getCalculatedField(device.getId(), device.getId());

56
dao/src/test/java/org/thingsboard/server/dao/service/validator/TenantProfileDataValidatorTest.java

@ -16,29 +16,35 @@
package org.thingsboard.server.dao.service.validator;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.dao.tenant.TenantProfileDao;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.exception.DataValidationException;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.Mockito.verify;
@SpringBootTest(classes = TenantProfileDataValidator.class)
class TenantProfileDataValidatorTest {
@MockBean
@MockitoBean
TenantProfileDao tenantProfileDao;
@MockBean
@MockitoBean
TenantProfileService tenantProfileService;
@SpyBean
@MockitoSpyBean
TenantProfileDataValidator validator;
TenantId tenantId = TenantId.fromUUID(UUID.fromString("9ef79cdf-37a8-4119-b682-2e7ed4e018da"));
@Test
@ -53,4 +59,44 @@ class TenantProfileDataValidatorTest {
verify(validator).validateString("Tenant profile name", tenantProfile.getName());
}
@ParameterizedTest
@ValueSource(ints = {-1, -100, Integer.MIN_VALUE})
void minAllowedScheduledUpdateIntervalInSecForCF_shouldRejectNegativeValues(int value) {
// GIVEN
var config = new DefaultTenantProfileConfiguration();
config.setMinAllowedScheduledUpdateIntervalInSecForCF(value);
var tenantProfileData = new TenantProfileData();
tenantProfileData.setConfiguration(config);
var tenantProfile = new TenantProfile();
tenantProfile.setName("Test");
tenantProfile.setProfileData(tenantProfileData);
// WHEN/THEN
assertThatThrownBy(() -> validator.validate(tenantProfile, __ -> TenantId.SYS_TENANT_ID))
.isInstanceOf(DataValidationException.class)
.hasMessageContaining("minAllowedScheduledUpdateIntervalInSecForCF")
.hasMessageContaining("must be greater than or equal to 0");
}
@ParameterizedTest
@ValueSource(ints = {0, 1, 60, Integer.MAX_VALUE})
void minAllowedScheduledUpdateIntervalInSecForCF_shouldAcceptValidValues(int value) {
// GIVEN
var config = new DefaultTenantProfileConfiguration();
config.setMinAllowedScheduledUpdateIntervalInSecForCF(value);
var tenantProfileData = new TenantProfileData();
tenantProfileData.setConfiguration(config);
var tenantProfile = new TenantProfile();
tenantProfile.setName("Test");
tenantProfile.setProfileData(tenantProfileData);
// WHEN/THEN
assertThatCode(() -> validator.validate(tenantProfile, __ -> TenantId.SYS_TENANT_ID))
.doesNotThrowAnyException();
}
}

2
edqs/src/main/resources/edqs.yml

@ -210,7 +210,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
health:
elasticsearch:
# Enable the org.springframework.boot.actuate.elasticsearch.ElasticsearchRestClientHealthIndicator.doHealthCheck

4
msa/vc-executor/src/main/resources/tb-vc-executor.yml

@ -215,6 +215,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -232,7 +234,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Service common properties
service:

43
rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java

@ -22,6 +22,7 @@ import com.google.common.base.Strings;
import lombok.Getter;
import lombok.SneakyThrows;
import org.apache.commons.io.IOUtils;
import org.apache.hc.core5.net.URIBuilder;
import org.apache.commons.lang3.concurrent.LazyInitializer;
import org.apache.hc.core5.net.URIBuilder;
import org.springframework.core.ParameterizedTypeReference;
@ -94,6 +95,8 @@ import org.thingsboard.server.common.data.asset.AssetSearchQuery;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.audit.AuditLog;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldInfo;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.common.data.domain.Domain;
import org.thingsboard.server.common.data.domain.DomainInfo;
@ -4368,6 +4371,46 @@ public class RestClient implements Closeable {
}
@SneakyThrows(URISyntaxException.class)
public PageData<CalculatedFieldInfo> getCalculatedFields(PageLink pageLink,
Set<CalculatedFieldType> types,
EntityType entityType,
Set<UUID> entities,
Set<String> names) {
var urlBuilder = new URIBuilder(baseURL).appendPath("/api/calculatedFields");
urlBuilder.addParameter("pageSize", String.valueOf(pageLink.getPageSize()));
urlBuilder.addParameter("page", String.valueOf(pageLink.getPage()));
if (!isEmpty(pageLink.getTextSearch())) {
urlBuilder.addParameter("textSearch", pageLink.getTextSearch());
}
if (pageLink.getSortOrder() != null) {
urlBuilder.addParameter("sortProperty", pageLink.getSortOrder().getProperty());
urlBuilder.addParameter("sortOrder", pageLink.getSortOrder().getDirection().name());
}
if (!CollectionUtils.isEmpty(types)) {
for (CalculatedFieldType type : types) {
urlBuilder.addParameter("types", type.name());
}
}
if (entityType != null) {
urlBuilder.addParameter("entityType", entityType.name());
}
if (!CollectionUtils.isEmpty(entities)) {
for (UUID entity : entities) {
urlBuilder.addParameter("entities", entity.toString());
}
}
if (!CollectionUtils.isEmpty(names)) {
for (String name : names) {
urlBuilder.addParameter("name", name);
}
}
return restTemplate.exchange(
urlBuilder.build(),
HttpMethod.GET, HttpEntity.EMPTY,
new ParameterizedTypeReference<PageData<CalculatedFieldInfo>>() {}).getBody();
}
public void deleteCalculatedField(CalculatedFieldId calculatedFieldId) {
restTemplate.delete(baseURL + "/api/calculatedField/{calculatedFieldId}", calculatedFieldId.getId());
}

31
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rest/TbRestApiCallNodeTest.java

@ -56,6 +56,7 @@ import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
@ExtendWith(MockitoExtension.class)
@ -115,18 +116,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
assertTrue(request.containsHeader("Foo"), "Custom header included");
assertEquals("Bar", request.getFirstHeader("Foo").getValue(), "Custom header value");
response.setStatusCode(200);
new Thread(new Runnable() {
@Override
public void run() {
try {
Thread.sleep(1000L);
} catch (InterruptedException e) {
// ignore
} finally {
latch.countDown();
}
}
}).start();
latch.countDown();
} catch (Exception e) {
System.out.println("Exception handling request: " + e.toString());
e.printStackTrace();
@ -158,7 +148,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
ArgumentCaptor<TbMsgMetaData> metadataCaptor = ArgumentCaptor.forClass(TbMsgMetaData.class);
ArgumentCaptor<String> dataCaptor = ArgumentCaptor.forClass(String.class);
verify(ctx).transformMsg(msgCaptor.capture(), metadataCaptor.capture(), dataCaptor.capture());
verify(ctx, timeout(10_000)).transformMsg(msgCaptor.capture(), metadataCaptor.capture(), dataCaptor.capture());
assertNotSame(metaData, metadataCaptor.getValue());
assertEquals(TbMsg.EMPTY_JSON_OBJECT, dataCaptor.getValue());
@ -184,18 +174,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
assertTrue(request.containsHeader("Foo"), "Custom header included");
assertEquals("Bar", request.getFirstHeader("Foo").getValue(), "Custom header value");
response.setStatusCode(200);
new Thread(new Runnable() {
@Override
public void run() {
try {
Thread.sleep(1000L);
} catch (InterruptedException e) {
// ignore
} finally {
latch.countDown();
}
}
}).start();
latch.countDown();
} catch (Exception e) {
System.out.println("Exception handling request: " + e.toString());
e.printStackTrace();
@ -227,7 +206,7 @@ public class TbRestApiCallNodeTest extends AbstractRuleNodeUpgradeTest {
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
ArgumentCaptor<TbMsgMetaData> metadataCaptor = ArgumentCaptor.forClass(TbMsgMetaData.class);
ArgumentCaptor<String> dataCaptor = ArgumentCaptor.forClass(String.class);
verify(ctx).transformMsg(msgCaptor.capture(), metadataCaptor.capture(), dataCaptor.capture());
verify(ctx, timeout(10_000)).transformMsg(msgCaptor.capture(), metadataCaptor.capture(), dataCaptor.capture());
assertNotSame(metaData, metadataCaptor.getValue());
assertEquals(TbMsg.EMPTY_JSON_OBJECT, dataCaptor.getValue());

4
transport/coap/src/main/resources/tb-coap-transport.yml

@ -421,6 +421,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -435,7 +437,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Notification system parameters
notification_system:

4
transport/http/src/main/resources/tb-http-transport.yml

@ -370,6 +370,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -384,7 +386,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Notification system parameters
notification_system:

4
transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml

@ -471,6 +471,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -485,7 +487,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Notification system parameters
notification_system:

4
transport/mqtt/src/main/resources/tb-mqtt-transport.yml

@ -404,6 +404,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -418,7 +420,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Notification system parameters
notification_system:

4
transport/snmp/src/main/resources/tb-snmp-transport.yml

@ -359,6 +359,8 @@ usage:
enabled_per_customer: "${USAGE_STATS_REPORT_PER_CUSTOMER_ENABLED:false}"
# Interval of reporting the statistics. By default, the summarized statistics are sent every 10 seconds
interval: "${USAGE_STATS_REPORT_INTERVAL:60}"
# Reporting interval for urgent keys (e.g. SMS, Email) that require quicker usage state updates
urgent_interval: "${USAGE_STATS_REPORT_URGENT_INTERVAL:10}"
# Amount of statistic messages in pack
pack_size: "${USAGE_STATS_REPORT_PACK_SIZE:1024}"
@ -373,7 +375,7 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: "${METRICS_ENDPOINTS_EXPOSE:info}"
# Notification system parameters
notification_system:

31
ui-ngx/patches/ngx-hm-carousel+19.0.0.patch

@ -0,0 +1,31 @@
diff --git a/node_modules/ngx-hm-carousel/fesm2022/ngx-hm-carousel.mjs b/node_modules/ngx-hm-carousel/fesm2022/ngx-hm-carousel.mjs
index 117f782..70e49a7 100644
--- a/node_modules/ngx-hm-carousel/fesm2022/ngx-hm-carousel.mjs
+++ b/node_modules/ngx-hm-carousel/fesm2022/ngx-hm-carousel.mjs
@@ -1,5 +1,5 @@
import * as i0 from '@angular/core';
-import { Directive, inject, ViewContainerRef, TemplateRef, input, PLATFORM_ID, DestroyRef, Renderer2, NgZone, ChangeDetectorRef, viewChild, ElementRef, contentChildren, contentChild, computed, signal, effect, forwardRef, Component, ChangeDetectionStrategy } from '@angular/core';
+import { Directive, inject, ViewContainerRef, TemplateRef, input, PLATFORM_ID, DestroyRef, Renderer2, NgZone, ChangeDetectorRef, viewChild, ElementRef, contentChildren, contentChild, computed, signal, effect, forwardRef, Component, ChangeDetectionStrategy, afterNextRender } from '@angular/core';
import { DOCUMENT, isPlatformBrowser, NgTemplateOutlet, AsyncPipe } from '@angular/common';
import { toObservable, takeUntilDestroyed } from '@angular/core/rxjs-interop';
import { NG_VALUE_ACCESSOR } from '@angular/forms';
@@ -340,7 +340,7 @@ class NgxHmCarouselComponent {
});
});
});
- const effectRef = effect(() => {
+ afterNextRender(() => {
this.rootElm = this.container().nativeElement;
this.containerElm = this.rootElm.children[0];
this.init();
@@ -365,10 +365,6 @@ class NgxHmCarouselComponent {
])
.pipe(takeUntilDestroyed(this._destroyRef))
.subscribe();
- // only exec once
- effectRef.destroy();
- }, {
- allowSignalWrites: true,
});
}
ngOnDestroy() {

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

@ -106,6 +106,8 @@ export class AlarmRulesTableConfig extends EntityTableConfig<AlarmRuleTableEntit
this.entityComponent = AlarmRulesComponent;
this.entityTabsComponent = AlarmRulesTabsComponent;
this.rowPointer = true;
} else {
this.addAsTextButton = false;
}
this.tableTitle = this.pageMode ? '' : this.translate.instant('alarm-rule.alarm-rules');
this.detailsPanelEnabled = this.pageMode;
@ -117,7 +119,7 @@ export class AlarmRulesTableConfig extends EntityTableConfig<AlarmRuleTableEntit
type: 'alarm-rule.alarm-rule',
typePlural: 'alarm-rule.alarm-rules',
list: 'alarm-rule.list',
add: 'action.add',
add: 'alarm-rule.add',
details: 'alarm-rule.details',
noEntities: 'alarm-rule.no-found',
search: 'action.search',

1
ui-ngx/src/app/modules/home/components/api-key/api-keys-table-config.ts

@ -59,7 +59,6 @@ export class ApiKeysTableConfig extends EntityTableConfig<ApiKeyInfo> {
this.entityType = EntityType.API_KEY;
this.detailsPanelEnabled = false;
this.addAsTextButton = true;
this.pageMode = false;
this.entityTranslations = entityTypeTranslations.get(EntityType.API_KEY);

2
ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.html

@ -142,7 +142,7 @@
<ng-container matColumnDef="select" sticky>
<mat-header-cell *matHeaderCellDef style="width: 40px;">
<mat-checkbox (change)="$event ? dataSource.masterToggle() : null"
[checked]="dataSource.selection.hasValue() && (dataSource.isAllSelected() | async)"
[(ngModel)]="selectAllModel"
[indeterminate]="dataSource.selection.hasValue() && !(dataSource.isAllSelected() | async)">
</mat-checkbox>
</mat-header-cell>

12
ui-ngx/src/app/modules/home/components/attribute/attribute-table.component.ts

@ -40,7 +40,7 @@ import { MatDialog } from '@angular/material/dialog';
import { DialogService } from '@core/services/dialog.service';
import { Direction, SortOrder } from '@shared/models/page/sort-order';
import { forkJoin, merge, Observable, Subject } from 'rxjs';
import { debounceTime, distinctUntilChanged, takeUntil } from 'rxjs/operators';
import { debounceTime, distinctUntilChanged, take, takeUntil } from 'rxjs/operators';
import { EntityId } from '@shared/models/id/entity-id';
import {
AttributeData,
@ -187,6 +187,7 @@ export class AttributeTableComponent extends PageComponent implements AfterViewI
textSearch = this.fb.control('', {nonNullable: true});
private destroy$ = new Subject<void>();
selectAllModel: boolean = false;
constructor(protected store: Store<AppState>,
private attributeService: AttributeService,
@ -223,6 +224,14 @@ export class AttributeTableComponent extends PageComponent implements AfterViewI
});
});
this.widgetResize$.observe(this.elementRef.nativeElement);
this.dataSource.selection.changed.pipe(
takeUntil(this.destroy$)
).subscribe(() => {
this.dataSource.isAllSelected().pipe(take(1)).subscribe(allSelected => {
this.selectAllModel = allSelected;
});
});
}
ngOnDestroy() {
@ -237,6 +246,7 @@ export class AttributeTableComponent extends PageComponent implements AfterViewI
this.attributeScope = attributeScope;
this.mode = 'default';
this.paginator.pageIndex = 0;
this.selectAllModel = false;
this.updateData(true);
}

2
ui-ngx/src/app/modules/home/components/calculated-fields/calculated-fields-table-config.ts

@ -110,6 +110,8 @@ export class CalculatedFieldsTableConfig extends EntityTableConfig<CalculatedFie
this.entityComponent = CalculatedFieldComponent;
this.entityTabsComponent = CalculatedFieldsTabsComponent;
this.rowPointer = true;
} else {
this.addAsTextButton = false;
}
this.tableTitle = this.translate.instant('entity.type-calculated-fields');
this.detailsPanelEnabled = this.pageMode;

104
ui-ngx/src/app/modules/home/components/entity/entities-table.component.html

@ -45,27 +45,19 @@
asButton strokedButton historyOnly [forAllTimeEnabled]="entitiesTableConfig.forAllTimeEnabled"></tb-timewindow>
</div>
<span class="flex-1"></span>
<div [class.!hidden]="!addEnabled()">
<ng-container *ngIf="!entitiesTableConfig.addActionDescriptors.length; else addActions">
<button *ngIf="!entitiesTableConfig.addAsTextButton"
mat-icon-button
<div [class.!hidden]="!addEnabled() || entitiesTableConfig.addAsTextButton">
<ng-container *ngIf="!entitiesTableConfig.addActionDescriptors.length; else iconCustom">
<button mat-icon-button
[disabled]="isLoading$ | async"
(click)="addEntity($event)"
matTooltip="{{ translations.add | translate }}"
matTooltipPosition="above">
<mat-icon>add</mat-icon>
</button>
<button *ngIf="entitiesTableConfig.addAsTextButton"
mat-stroked-button color="primary"
[disabled]="isLoading$ | async"
(click)="addEntity($event)">
<mat-icon>add</mat-icon>
{{ translations.add | translate }}
</button>
</ng-container>
<ng-template #addActions>
<ng-container *ngIf="entitiesTableConfig.addActionDescriptors.length === 1; else addActionsMenu">
<button *ngIf="!entitiesTableConfig.addAsTextButton && entitiesTableConfig.addActionDescriptors[0].isEnabled()"
<ng-template #iconCustom>
<ng-container *ngIf="entitiesTableConfig.addActionDescriptors.length === 1; else iconMenu">
<button *ngIf="entitiesTableConfig.addActionDescriptors[0].isEnabled()"
mat-icon-button
[disabled]="isLoading$ | async"
(click)="entitiesTableConfig.addActionDescriptors[0].onAction($event)"
@ -73,30 +65,14 @@
matTooltipPosition="above">
<tb-icon>{{entitiesTableConfig.addActionDescriptors[0].icon}}</tb-icon>
</button>
<button *ngIf="entitiesTableConfig.addAsTextButton && entitiesTableConfig.addActionDescriptors[0].isEnabled()"
mat-stroked-button color="primary"
[disabled]="isLoading$ | async"
(click)="entitiesTableConfig.addActionDescriptors[0].onAction($event)">
<tb-icon>{{entitiesTableConfig.addActionDescriptors[0].icon}}</tb-icon>
{{ entitiesTableConfig.addActionDescriptors[0].name }}
</button>
</ng-container>
<ng-template #addActionsMenu>
<ng-template #iconMenu>
<button mat-icon-button [disabled]="isLoading$ | async"
matTooltip="{{ translations.add | translate }}"
matTooltipPosition="above"
[matMenuTriggerFor]="addActionsMenu">
<mat-icon>add</mat-icon>
</button>
<mat-menu #addActionsMenu="matMenu" xPosition="before">
<button mat-menu-item *ngFor="let actionDescriptor of entitiesTableConfig.addActionDescriptors"
[disabled]="isLoading$ | async"
[class.!hidden]="!actionDescriptor.isEnabled()"
(click)="actionDescriptor.onAction($event)">
<tb-icon matMenuItemIcon>{{actionDescriptor.icon}}</tb-icon>
<span>{{ actionDescriptor.name }}</span>
</button>
</mat-menu>
</ng-template>
</ng-template>
</div>
@ -118,6 +94,72 @@
matTooltipPosition="above">
<mat-icon>search</mat-icon>
</button>
<div [class.!hidden]="!addEnabled() || !entitiesTableConfig.addAsTextButton" class="pl-3" style="--mat-fab-small-container-elevation-shadow: none;">
<ng-container *ngIf="!entitiesTableConfig.addActionDescriptors.length; else textActions">
<button class="lt-sm:!hidden"
mat-flat-button color="primary"
[disabled]="isLoading$ | async"
(click)="addEntity($event)">
<mat-icon>add</mat-icon>
{{ translations.add | translate }}
</button>
<button class="gt-xs:!hidden"
mat-mini-fab color="primary"
[disabled]="isLoading$ | async"
(click)="addEntity($event)"
matTooltip="{{ translations.add | translate }}"
matTooltipPosition="above">
<mat-icon>add</mat-icon>
</button>
</ng-container>
<ng-template #textActions>
<ng-container *ngIf="entitiesTableConfig.addActionDescriptors.length === 1; else textMenu">
<ng-container *ngIf="entitiesTableConfig.addActionDescriptors[0].isEnabled()">
<button class="lt-sm:!hidden"
mat-flat-button color="primary"
[disabled]="isLoading$ | async"
(click)="entitiesTableConfig.addActionDescriptors[0].onAction($event)">
<tb-icon matButtonIcon>{{entitiesTableConfig.addActionDescriptors[0].icon}}</tb-icon>
{{ entitiesTableConfig.addActionDescriptors[0].name }}
</button>
<button class="gt-xs:!hidden"
mat-mini-fab color="primary"
[disabled]="isLoading$ | async"
(click)="entitiesTableConfig.addActionDescriptors[0].onAction($event)"
matTooltip="{{ entitiesTableConfig.addActionDescriptors[0].name }}"
matTooltipPosition="above">
<mat-icon>{{entitiesTableConfig.addActionDescriptors[0].icon}}</mat-icon>
</button>
</ng-container>
</ng-container>
<ng-template #textMenu>
<button class="lt-sm:!hidden"
mat-flat-button color="primary"
[disabled]="isLoading$ | async"
[matMenuTriggerFor]="addActionsMenu">
<mat-icon>add</mat-icon>
{{ translations.add | translate }}
</button>
<button class="gt-xs:!hidden"
mat-mini-fab color="primary"
[disabled]="isLoading$ | async"
[matMenuTriggerFor]="addActionsMenu"
matTooltip="{{ translations.add | translate }}"
matTooltipPosition="above">
<mat-icon>add</mat-icon>
</button>
</ng-template>
</ng-template>
</div>
<mat-menu #addActionsMenu="matMenu" xPosition="before">
<button mat-menu-item *ngFor="let actionDescriptor of entitiesTableConfig.addActionDescriptors"
[disabled]="isLoading$ | async"
[class.!hidden]="!actionDescriptor.isEnabled()"
(click)="actionDescriptor.onAction($event)">
<tb-icon matMenuItemIcon>{{actionDescriptor.icon}}</tb-icon>
<span>{{ actionDescriptor.name }}</span>
</button>
</mat-menu>
</div>
</mat-toolbar>
<mat-toolbar class="mat-mdc-table-toolbar" [class.!hidden]="!textSearchMode || !dataSource.selection.isEmpty()">

2
ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.scss

@ -43,7 +43,7 @@
}
}
> .mat-expansion-panel-content {
.mat-expansion-panel-content {
> .mat-expansion-panel-body {
padding: 0;
}

18
ui-ngx/src/app/modules/home/components/vc/entity-versions-table.component.html

@ -32,7 +32,8 @@
</tb-branch-autocomplete>
</div>
<span class="flex-1"></span>
<div class="flex flex-row items-center justify-center xs:flex-col xs:items-end xs:justify-center">
<div class="flex flex-row items-center justify-center xs:flex-col xs:items-end xs:justify-center"
[class.xs:flex-col-reverse]="!singleEntityMode">
<button *ngIf="singleEntityMode" mat-stroked-button color="primary"
#createVersionButton
[disabled]="(isLoading$ | async) || (isReadOnly | async)"
@ -40,13 +41,6 @@
<mat-icon>update</mat-icon>
{{'version-control.create-version' | translate }}
</button>
<button *ngIf="!singleEntityMode" mat-stroked-button color="primary"
#complexCreateVersionButton
[disabled]="(isLoading$ | async) || (isReadOnly | async)"
(click)="toggleComplexCreateVersion($event, complexCreateVersionButton)">
<mat-icon>update</mat-icon>
{{'version-control.create-entities-version' | translate }}
</button>
<div class="flex flex-row">
<button mat-icon-button [disabled]="isLoading$ | async" (click)="updateData()"
matTooltip="{{ 'action.refresh' | translate }}"
@ -61,6 +55,14 @@
<mat-icon>search</mat-icon>
</button>
</div>
<button *ngIf="!singleEntityMode" mat-flat-button color="primary"
class="gt-xs:ml-3"
#complexCreateVersionButton
[disabled]="(isLoading$ | async) || (isReadOnly | async)"
(click)="toggleComplexCreateVersion($event, complexCreateVersionButton)">
<mat-icon>update</mat-icon>
{{'version-control.create-entities-version' | translate }}
</button>
</div>
</div>
</mat-toolbar>

11
ui-ngx/src/app/modules/home/components/widget/lib/home-page/recent-dashboards-widget.component.html

@ -24,8 +24,15 @@
<tb-toggle-option value="last">{{ 'widgets.recent-dashboards.last' | translate }}</tb-toggle-option>
<tb-toggle-option value="starred">{{ 'widgets.recent-dashboards.starred' | translate }}</tb-toggle-option>
</tb-toggle-header>
<a *ngIf="authUser.authority === authority.TENANT_ADMIN" class="md:!hidden"
mat-flat-button color="primary" routerLink="/dashboards" [queryParams]="{action: 'add'}">{{ 'dashboard.add' | translate }}</a>
<ng-container *ngIf="authUser.authority === authority.TENANT_ADMIN">
@if (hasDevice) {
<a class="md:!hidden"
mat-flat-button color="primary" routerLink="/dashboards/all" [queryParams]="{action: 'add'}">{{ 'dashboard.add' | translate }}</a>
} @else {
<a class="md:!hidden"
mat-stroked-button color="primary" routerLink="/dashboards/all" [queryParams]="{action: 'add'}">{{ 'dashboard.add' | translate }}</a>
}
</ng-container>
</div>
</div>
<ng-container *ngIf="userDashboardsInfo; else loading" [ngSwitch]="dashboardsToggle.value">

45
ui-ngx/src/app/modules/home/components/widget/lib/home-page/recent-dashboards-widget.component.ts

@ -48,6 +48,13 @@ import { Direction, SortOrder } from '@shared/models/page/sort-order';
import { MatSort } from '@angular/material/sort';
import { DashboardInfo } from '@shared/models/dashboard.models';
import { DashboardAutocompleteComponent } from '@shared/components/dashboard-autocomplete.component';
import { UtilsService } from '@core/services/utils.service';
import { Datasource, DatasourceType, widgetType } from '@shared/models/widget.models';
import { IWidgetSubscription, WidgetSubscriptionOptions } from '@core/api/widget-api.models';
import { formattedDataFormDatasourceData } from '@core/utils';
import { AliasFilterType } from '@shared/models/alias.models';
import { DataKeyType } from '@shared/models/telemetry/telemetry.models';
import { EntityType } from '@shared/models/entity-type.models';
@Component({
selector: 'tb-recent-dashboards-widget',
@ -78,13 +85,16 @@ export class RecentDashboardsWidgetComponent extends PageComponent implements On
starredDashboardValue = null;
hasDashboardsAccess = true;
hasDevice = true;
dirty = false;
public customerId: string;
private isFullscreenMode = getCurrentAuthState(this.store).forceFullscreen;
private subscription: IWidgetSubscription;
constructor(protected store: Store<AppState>,
private cd: ChangeDetectorRef,
private utils: UtilsService,
private userSettingService: UserSettingsService) {
super(store);
}
@ -96,6 +106,41 @@ export class RecentDashboardsWidgetComponent extends PageComponent implements On
this.hasDashboardsAccess = [Authority.TENANT_ADMIN, Authority.CUSTOMER_USER].includes(this.authUser.authority);
if (this.hasDashboardsAccess) {
this.reload();
if (window.location.pathname.startsWith('/home') && this.authUser.authority === Authority.TENANT_ADMIN) {
const ds: Datasource = {
type: DatasourceType.entityCount,
name: '',
entityFilter: {
entityType: EntityType.DEVICE,
type: AliasFilterType.entityType
},
dataKeys: [this.utils.createKey({ name: 'count'}, DataKeyType.count)]
}
const apiUsageSubscriptionOptions: WidgetSubscriptionOptions = {
datasources: [ds],
useDashboardTimewindow: false,
type: widgetType.latest,
callbacks: {
onDataUpdated: (subscription) => {
const data = formattedDataFormDatasourceData(subscription.data);
this.hasDevice = (data[0].count || 0) !== 0;
this.cd.detectChanges();
}
}
};
this.ctx.subscriptionApi.createSubscription(apiUsageSubscriptionOptions, true).subscribe((subscription) => {
this.subscription = subscription;
});
}
}
}
ngOnDestroy() {
super.ngOnDestroy();
if (this.subscription) {
this.ctx.subscriptionApi.removeSubscription(this.subscription.id);
}
}

40
ui-ngx/src/app/modules/home/components/widget/lib/rpc/power-button-widget.models.ts

@ -337,6 +337,9 @@ export abstract class PowerButtonShape {
protected onPowerSymbolLine: Path;
private onIcon$: Observable<Element>;
private offIcon$: Observable<Element>;
protected onIconOffsetX: number = 0;
protected onIconOffsetY: number = 0;
protected constructor(protected widgetContext: WidgetContext,
protected svgShape: Svg,
@ -362,10 +365,14 @@ export abstract class PowerButtonShape {
take(1),
map((svgElement) => {
const element = new Element(svgElement.firstChild);
const iconGroup = this.svgShape.group();
element.addTo(iconGroup);
const box = element.bbox();
const scale = size / box.height;
element.scale(scale);
return element;
const scaledBox = iconGroup.bbox();
iconGroup.translate(-scaledBox.cx, -scaledBox.cy);
return iconGroup;
}),
catchError(() => of(null)
));
@ -416,8 +423,15 @@ export abstract class PowerButtonShape {
public drawOffShape(centerGroup: G, label: boolean, labelWeight?: string, circleStroke?: boolean) {
if (this.icons.offButtonIcon.showIcon) {
this.createIconElement(this.icons.offButtonIcon.icon, this.icons.offButtonIcon.iconSize).subscribe(icon =>
this.offPowerSymbolIcon = icon.center(cx, cy).addTo(centerGroup));
this.offIcon$ = this.createIconElement(this.icons.offButtonIcon.icon, this.icons.offButtonIcon.iconSize)
.pipe(shareReplay(1));
this.offIcon$.pipe(take(1)).subscribe(icon => {
this.offPowerSymbolIcon = icon;
icon.translate(cx, cy);
centerGroup.add(icon);
});
} else {
if (label) {
this.offLabelShape = this.createOffLabel(labelWeight).addTo(centerGroup);
@ -434,9 +448,12 @@ export abstract class PowerButtonShape {
.pipe(shareReplay(1));
this.onIcon$.subscribe(icon => {
this.onPowerSymbolIcon = icon.center(cx, cy);
const iconBox = icon.bbox();
this.onIconOffsetX = iconBox.cx;
this.onIconOffsetY = iconBox.cy;
this.onPowerSymbolIcon = icon.translate(cx, cy);
if (isDefinedAndNotNull(onCenterGroup)) {
this.onPowerSymbolIcon.addTo(onCenterGroup);
onCenterGroup.add(this.onPowerSymbolIcon);
}
if (isDefinedAndNotNull(mask)) {
this.createMask(mask, [this.onPowerSymbolIcon]);
@ -468,7 +485,7 @@ export abstract class PowerButtonShape {
public onCenterTimeLine(timeline: Timeline, label: boolean) {
if (this.icons.onButtonIcon.showIcon) {
if (this.onIcon$) {
this.onIcon$.subscribe(icon => icon.timeline(timeline))
this.onIcon$.subscribe(icon => icon.timeline(timeline));
}
} else {
if (label) {
@ -482,7 +499,7 @@ export abstract class PowerButtonShape {
public offCenterColor(mainColor: PowerButtonColor, label: boolean) {
if (this.icons.offButtonIcon.showIcon) {
this.offPowerSymbolIcon.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity});
this.offIcon$.subscribe(icon => icon.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity}))
} else {
if (label) {
this.offLabelShape.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity});
@ -495,7 +512,7 @@ export abstract class PowerButtonShape {
public onCenterColor(mainColor: PowerButtonColor, label: boolean) {
if (this.icons.onButtonIcon.showIcon) {
this.onPowerSymbolIcon.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity});
this.onIcon$.subscribe((icon)=> icon.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity}))
} else {
if (label) {
this.onLabelShape.attr({ fill: mainColor.hex, 'fill-opacity': mainColor.opacity});
@ -508,7 +525,12 @@ export abstract class PowerButtonShape {
public buttonAnimation(scale: number, label: boolean) {
if (this.icons.onButtonIcon.showIcon) {
powerButtonAnimation(this.onPowerSymbolIcon).transform({scale});
const translateX = cx - this.onIconOffsetX;
const translateY = cy - this.onIconOffsetY;
powerButtonAnimation(this.onPowerSymbolIcon).transform({
scale: scale,
translate: [translateX, translateY]
});
} else {
if (label) {
powerButtonAnimation(this.onLabelShape).transform({scale, origin: {x: cx, y: cy}});

2
ui-ngx/src/app/modules/home/components/widget/lib/settings/common/image-cards-select.component.scss

@ -44,7 +44,7 @@
display: none;
}
}
> .mat-expansion-panel-content {
.mat-expansion-panel-content {
> .mat-expansion-panel-body {
padding: 0 16px 16px !important;
}

5
ui-ngx/src/app/modules/home/components/widget/lib/settings/common/map/marker-image-settings.component.ts

@ -82,11 +82,12 @@ export class MarkerImageSettingsComponent implements ControlValueAccessor {
renderer: this.renderer,
componentType: MarkerImageSettingsPanelComponent,
hostView: this.viewContainerRef,
preferredPlacement: 'left',
preferredPlacement: 'leftOnly',
context: {
markerImageSettings: this.modelValue,
},
isModal: true
isModal: true,
overlayStyle: {padding: '10px'}
});
markerImageSettingsPanelPopover.tbComponentRef.instance.popover = markerImageSettingsPanelPopover;
markerImageSettingsPanelPopover.tbComponentRef.instance.markerImageSettingsApplied.subscribe((markerImageSettings) => {

5
ui-ngx/src/app/modules/home/components/widget/lib/settings/common/map/shape-fill-image-settings.component.ts

@ -82,11 +82,12 @@ export class ShapeFillImageSettingsComponent implements ControlValueAccessor {
renderer: this.renderer,
componentType: ShapeFillImageSettingsPanelComponent,
hostView: this.viewContainerRef,
preferredPlacement: 'left',
preferredPlacement: 'leftOnly',
context: {
shapeFillImageSettings: this.modelValue,
},
isModal: true
isModal: true,
overlayStyle: {padding: '10px'}
}).tbComponentRef.instance.shapeFillImageSettingsApplied.subscribe((shapeFillImageSettings) => {
this.modelValue = shapeFillImageSettings;
this.propagateChange(this.modelValue);

2
ui-ngx/src/app/modules/home/components/widget/lib/settings/widget-settings.scss

@ -108,7 +108,7 @@
align-items: center;
}
> .mat-expansion-panel-content {
.mat-expansion-panel-content {
> .mat-expansion-panel-body {
padding: 0;
}

6
ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts

@ -115,9 +115,11 @@ export class AttributeDatasource implements DataSource<AttributeData> {
}
isAllSelected(): Observable<boolean> {
const numSelected = this.selection.selected.length;
return this.attributesSubject.pipe(
map((attributes) => numSelected === attributes.length)
map((attributes) => {
const numSelected = this.selection.selected.length;
return attributes.length > 0 && numSelected === attributes.length
})
);
}

2
ui-ngx/src/app/modules/home/models/entity/entities-table-config.models.ts

@ -175,7 +175,7 @@ export class EntityTableConfig<T extends BaseData<HasId>, P extends PageLink = P
selectionEnabled = true;
searchEnabled = true;
addEnabled = true;
addAsTextButton = false;
addAsTextButton = true;
entitiesDeleteEnabled = true;
detailsPanelEnabled = true;
hideDetailsTabsOnEdit = true;

2
ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html

@ -280,7 +280,7 @@
</mat-error>
</div>
<div class="flex flex-1 flex-col">
<div class="flex flex-1 flex-col gt-xs:w-0">
<mat-form-field class="mat-block flex-1">
<mat-label translate>admin.oauth2.redirect-uri-template</mat-label>
<input matInput formControlName="redirectUri" readonly>

1
ui-ngx/src/app/modules/home/pages/ai-model/ai-model-table-config.resolve.ts

@ -47,7 +47,6 @@ export class AiModelsTableConfigResolver {
) {
this.config.selectionEnabled = true;
this.config.entityType = EntityType.AI_MODEL;
this.config.addAsTextButton = true;
this.config.rowPointer = true;
this.config.detailsPanelEnabled = false;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.AI_MODEL);

1
ui-ngx/src/app/modules/home/pages/mobile/applications/mobile-app-table-config.resolver.ts

@ -57,7 +57,6 @@ export class MobileAppTableConfigResolver {
) {
this.config.selectionEnabled = false;
this.config.entityType = EntityType.MOBILE_APP;
this.config.addAsTextButton = true;
this.config.entitiesDeleteEnabled = false;
this.config.rowPointer = true;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.MOBILE_APP);

1
ui-ngx/src/app/modules/home/pages/mobile/bundes/mobile-bundle-table-config.resolve.ts

@ -63,7 +63,6 @@ export class MobileBundleTableConfigResolver {
) {
this.config.selectionEnabled = false;
this.config.entityType = EntityType.MOBILE_APP_BUNDLE;
this.config.addAsTextButton = true;
this.config.rowPointer = true;
this.config.detailsPanelEnabled = false;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.MOBILE_APP_BUNDLE);

2
ui-ngx/src/app/modules/home/pages/mobile/common/editor-panel.component.ts

@ -48,7 +48,7 @@ export class EditorPanelComponent implements OnInit {
tinyMceOptions: Partial<EditorOptions> = {
base_url: '/assets/tinymce',
suffix: '.min',
plugins: ['link', 'table', 'image', 'imagetools', 'lists', 'fullscreen'],
plugins: ['link', 'table', 'image', 'lists', 'fullscreen'],
menubar: 'edit insert view format',
toolbar: ['fontfamily fontsize | bold italic underline strikethrough forecolor backcolor',
'alignleft aligncenter alignright alignjustify | bullist | link table image | fullscreen'],

1
ui-ngx/src/app/modules/home/pages/notification/recipient/recipient-table-config.resolver.ts

@ -48,7 +48,6 @@ export class RecipientTableConfigResolver {
this.config.entityType = EntityType.NOTIFICATION_TARGET;
this.config.detailsPanelEnabled = false;
this.config.addAsTextButton = true;
this.config.rowPointer = true;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.NOTIFICATION_TARGET);

1
ui-ngx/src/app/modules/home/pages/notification/rule/rule-table-config.resolver.ts

@ -49,7 +49,6 @@ export class RuleTableConfigResolver {
this.config.entityType = EntityType.NOTIFICATION_RULE;
this.config.detailsPanelEnabled = false;
this.config.addAsTextButton = true;
this.config.rowPointer = true;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.NOTIFICATION_RULE);

2
ui-ngx/src/app/modules/home/pages/notification/template/configuration/notification-template-configuration.component.scss

@ -65,7 +65,7 @@
:host ::ng-deep {
.tb-form-panel .mat-expansion-panel.tb-settings {
& > .mat-expansion-panel-content > .mat-expansion-panel-body {
.mat-expansion-panel-content > .mat-expansion-panel-body {
gap: 0;
}
}

1
ui-ngx/src/app/modules/home/pages/notification/template/template-table-config.resolver.ts

@ -47,7 +47,6 @@ export class TemplateTableConfigResolver {
this.config.entityType = EntityType.NOTIFICATION_TEMPLATE;
this.config.detailsPanelEnabled = false;
this.config.addAsTextButton = true;
this.config.rowPointer = true;
this.config.entityTranslations = entityTypeTranslations.get(EntityType.NOTIFICATION_TEMPLATE);

2
ui-ngx/src/app/modules/login/pages/login/login.component.ts

@ -76,7 +76,7 @@ export class LoginComponent extends PageComponent implements OnInit {
getOAuth2Uri(oauth2Client: OAuth2ClientLoginInfo): string {
let result = "";
if (this.authService.redirectUrl) {
result += "?prevUri=" + this.authService.redirectUrl;
result += "?prevUri=" + encodeURIComponent(this.authService.redirectUrl);
}
return oauth2Client.url + result;
}

7
ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts

@ -14,9 +14,9 @@
/// limitations under the License.
///
import { OverlayModule } from '@angular/cdk/overlay';
import { Overlay, OverlayContainer, OverlayModule } from '@angular/cdk/overlay';
import { NgModule } from '@angular/core';
import { DEFAULT_DIALOG_CONFIG, DialogConfig, DialogModule } from '@angular/cdk/dialog';
import { DEFAULT_DIALOG_CONFIG, Dialog, DialogConfig, DialogModule } from '@angular/cdk/dialog';
import { MatDialogModule } from '@angular/material/dialog';
import { DynamicDialog, DynamicMatDialog } from './dynamic-dialog';
import { DynamicOverlay } from './dynamic-overlay';
@ -24,8 +24,11 @@ import { DynamicOverlayContainer } from './dynamic-overlay-container';
export const DYNAMIC_MAT_DIALOG_PROVIDERS = [
DynamicOverlayContainer,
{ provide: OverlayContainer, useExisting: DynamicOverlayContainer },
DynamicOverlay,
{ provide: Overlay, useExisting: DynamicOverlay },
DynamicDialog,
{ provide: Dialog, useExisting: DynamicDialog },
DynamicMatDialog,
{
provide: DEFAULT_DIALOG_CONFIG,

45
ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts

@ -14,18 +14,11 @@
/// limitations under the License.
///
import { Location } from "@angular/common";
import { Inject, Injectable, Injector, Optional, SkipSelf, TemplateRef } from '@angular/core';
import {
MAT_DIALOG_DEFAULT_OPTIONS,
MAT_DIALOG_SCROLL_STRATEGY,
MatDialog,
MatDialogConfig, MatDialogRef
} from '@angular/material/dialog';
import { DynamicOverlay } from "./dynamic-overlay";
import { DynamicOverlayContainer } from '@shared/components/dialog/dynamic/dynamic-overlay-container';
import { ComponentType, ScrollStrategy } from '@angular/cdk/overlay';
import { DEFAULT_DIALOG_CONFIG, Dialog, DialogConfig } from '@angular/cdk/dialog';
import { inject, Injectable, TemplateRef } from '@angular/core';
import { MatDialog, MatDialogConfig, MatDialogRef } from '@angular/material/dialog';
import { DynamicOverlay } from './dynamic-overlay';
import { ComponentType } from '@angular/cdk/overlay';
import { Dialog } from '@angular/cdk/dialog';
export interface DynamicMatDialogConfig<D> extends MatDialogConfig<D> {
containerElement?: HTMLElement;
@ -34,25 +27,12 @@ export interface DynamicMatDialogConfig<D> extends MatDialogConfig<D> {
@Injectable()
export class DynamicMatDialog extends MatDialog {
private _customOverlay: DynamicOverlay;
private _customOverlay = inject(DynamicOverlay);
constructor( _overlay: DynamicOverlay,
_injector: Injector,
@Optional() location: Location,
@Inject( MAT_DIALOG_DEFAULT_OPTIONS ) _defaultOptions: MatDialogConfig,
@Inject( MAT_DIALOG_SCROLL_STRATEGY ) _scrollStrategy: ScrollStrategy,
@Optional() @SkipSelf() _parentDialog:DynamicMatDialog,
_overlayContainer: DynamicOverlayContainer) {
super( _overlay, _injector, location, _defaultOptions, _scrollStrategy, _parentDialog, _overlayContainer );
this._dialog = _injector.get(DynamicDialog);
this._customOverlay = _overlay;
}
public open<T, D = any, R = any>(component: ComponentType<T> | TemplateRef<T>, config?: DynamicMatDialogConfig<D>): MatDialogRef<T, R> {
public override open<T, D = any, R = any>(component: ComponentType<T> | TemplateRef<T>, config?: DynamicMatDialogConfig<D>): MatDialogRef<T, R> {
if (config?.containerElement) {
config.containerElement.style.transform = 'translateZ(0)';
this._customOverlay.setContainerElement( config.containerElement );
this._customOverlay.setContainerElement(config.containerElement);
}
const ref = super.open(component, config);
if (config?.containerElement) {
@ -73,13 +53,4 @@ export class DynamicMatDialog extends MatDialog {
@Injectable()
export class DynamicDialog extends Dialog {
constructor( _overlay: DynamicOverlay,
_injector: Injector,
@Inject( DEFAULT_DIALOG_CONFIG ) _defaultOptions: DialogConfig,
@Inject( MAT_DIALOG_SCROLL_STRATEGY ) _scrollStrategy: ScrollStrategy,
@Optional() @SkipSelf() _parentDialog: DynamicDialog,
_overlayContainer: DynamicOverlayContainer) {
super( _overlay, _injector, _defaultOptions, _parentDialog, _overlayContainer, _scrollStrategy );
}
}

43
ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay.ts

@ -14,49 +14,16 @@
/// limitations under the License.
///
import {
Overlay,
ScrollStrategyOptions,
OverlayKeyboardDispatcher, OverlayOutsideClickDispatcher, OverlayPositionBuilder
} from '@angular/cdk/overlay';
import { ComponentFactoryResolver, Inject, Injectable, Injector, NgZone, DOCUMENT } from '@angular/core';
import { Overlay } from '@angular/cdk/overlay';
import { inject, Injectable } from '@angular/core';
import { DynamicOverlayContainer } from './dynamic-overlay-container';
import { Location } from '@angular/common';
import { Directionality } from '@angular/cdk/bidi';
@Injectable()
export class DynamicOverlay extends Overlay {
private _dynamicOverlayContainer: DynamicOverlayContainer;
private _dynamicOverlayContainer = inject(DynamicOverlayContainer);
constructor( scrollStrategies: ScrollStrategyOptions,
_overlayContainer: DynamicOverlayContainer,
_componentFactoryResolver: ComponentFactoryResolver,
_positionBuilder: OverlayPositionBuilder,
_keyboardDispatcher: OverlayKeyboardDispatcher,
_injector: Injector,
_ngZone: NgZone,
@Inject(DOCUMENT) document: Document,
_directionality: Directionality,
_location: Location,
_outsideClickDispatcher: OverlayOutsideClickDispatcher) {
super( scrollStrategies,
_overlayContainer,
_componentFactoryResolver,
_positionBuilder,
_keyboardDispatcher,
_injector,
_ngZone,
document,
_directionality,
_location,
_outsideClickDispatcher);
this._dynamicOverlayContainer = _overlayContainer;
}
public setContainerElement(containerElement:HTMLElement ): void {
this._dynamicOverlayContainer.setContainerElement( containerElement );
public setContainerElement(containerElement: HTMLElement): void {
this._dynamicOverlayContainer.setContainerElement(containerElement);
}
}

6
ui-ngx/src/app/shared/components/image/image-gallery.component.html

@ -74,10 +74,10 @@
<tb-icon>mdi:file-import</tb-icon>
</button>
<button [class.!hidden]="dialogMode && isScada"
class="gt-md:!hidden"
class="ml-3 gt-md:!hidden"
[disabled]="isLoading$ | async"
mat-icon-button
color="primary"
mat-mini-fab color="primary"
style="--mat-fab-small-container-elevation-shadow: none;"
(click)="uploadImage()"
matTooltip="{{ (isScada ? 'scada.upload-symbol' : 'image.upload-image' ) | translate }}"
matTooltipPosition="above">

1
ui-ngx/src/app/shared/models/entity-type.models.ts

@ -490,6 +490,7 @@ export const entityTypeTranslations = new Map<EntityType | AliasEntityType, Enti
type: 'entity.type-calculated-field',
typePlural: 'entity.type-calculated-fields',
list: 'calculated-fields.list',
add: 'calculated-fields.add',
details: 'calculated-fields.calculated-field-details',
noEntities: 'calculated-fields.no-found',
search: 'action.search',

2
ui-ngx/src/assets/dashboard/tenant_admin_home_page.json

@ -224,7 +224,7 @@
"padding": "16px",
"settings": {
"useMarkdownTextFunction": false,
"markdownTextPattern": "<div class=\"tb-card-content\">\n <div class=\"tb-content-container\">\n <div class=\"tb-card-header\">\n <div class=\"tb-card-title\">\n <a class=\"tb-home-widget-title tb-home-widget-link\" routerLink=\"/entities/devices\">{{ 'device.devices' | translate }}</a>\n </div>\n <div class=\"flex flex-row gap-3\">\n <a class=\"md:!hidden\" mat-stroked-button color=\"primary\" href=\"https://thingsboard.io/docs\" target=\"_blank\">{{ 'widgets.devices.view-docs' | translate }}</a>\n <a mat-flat-button color=\"primary\" routerLink=\"/entities/devices\" [queryParams]=\"{action: 'add'}\">{{ 'device.add' | translate }}</a>\n </div>\n </div>\n <div class=\"tb-item-cards\">\n <a class=\"tb-item-card tb-inactive\" routerLink=\"/entities/devices\" [queryParams]=\"{active: false}\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.inactive</div>\n </div>\n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${inactiveDevices:0}</div>\n </div>\n </a>\n <a class=\"tb-item-card tb-active\" routerLink=\"/entities/devices\" [queryParams]=\"{active: true}\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.active</div>\n </div> \n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${activeDevices:0}</div>\n </div>\n </a>\n <a class=\"tb-item-card tb-total md:!hidden\" routerLink=\"/entities/devices\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.total</div>\n </div> \n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${totalDevices:0}</div>\n </div>\n </a>\n </div>\n </div>\n</div>",
"markdownTextPattern": "<div class=\"tb-card-content\">\n <div class=\"tb-content-container\">\n <div class=\"tb-card-header\">\n <div class=\"tb-card-title\">\n <a class=\"tb-home-widget-title tb-home-widget-link\" routerLink=\"/entities/devices\">{{ 'device.devices' | translate }}</a>\n </div>\n <div class=\"flex flex-row gap-3\">\n <a class=\"md:!hidden\" mat-button color=\"primary\" href=\"https://thingsboard.io/docs\" target=\"_blank\">{{ 'widgets.devices.view-docs' | translate }}</a>\n <a *ngIf=\"${totalDevices:0} === 0\" mat-flat-button color=\"primary\" routerLink=\"/entities/devices\" [queryParams]=\"{action: 'add'}\">{{ 'device.add' | translate }}</a>\n <a *ngIf=\"${totalDevices:0} !== 0\" mat-stroked-button color=\"primary\" routerLink=\"/entities/devices\" [queryParams]=\"{action: 'add'}\">{{ 'device.add' | translate }}</a>\n </div>\n </div>\n <div class=\"tb-item-cards\">\n <a class=\"tb-item-card tb-inactive\" routerLink=\"/entities/devices\" [queryParams]=\"{active: false}\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.inactive</div>\n </div>\n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${inactiveDevices:0}</div>\n </div>\n </a>\n <a class=\"tb-item-card tb-active\" routerLink=\"/entities/devices\" [queryParams]=\"{active: true}\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.active</div>\n </div> \n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${activeDevices:0}</div>\n </div>\n </a>\n <a class=\"tb-item-card tb-total md:!hidden\" routerLink=\"/entities/devices\">\n <div class=\"tb-item-title-container\">\n <div class=\"tb-item-title tb-home-widget-link\" translate>widgets.devices.total</div>\n </div> \n <div class=\"tb-count-container\">\n <div class=\"tb-count\">${totalDevices:0}</div>\n </div>\n </a>\n </div>\n </div>\n</div>",
"applyDefaultMarkdownStyle": false,
"markdownCss": ".tb-card-content {\n width: 100%;\n height: 100%;\n display: flex;\n flex-direction: row;\n}\n\n.tb-content-container {\n flex: 1;\n display: flex;\n flex-direction: column;\n justify-content: flex-start;\n gap: 12px;\n}\n\n.tb-card-header {\n height: 36px;\n display: flex;\n flex-direction: row;\n justify-content: space-between;\n}\n\n.tb-item-cards {\n flex: 1;\n display: flex;\n flex-direction: row;\n gap: 12px;\n overflow: hidden;\n}\n\na.tb-item-card {\n flex: 1;\n display: flex;\n flex-direction: column;\n padding: 8px 12px;\n border: 1px solid;\n border-radius: 10px;\n margin-bottom: 12px;\n overflow: hidden;\n justify-content: space-evenly;\n}\n\na.tb-item-card.tb-inactive {\n background: rgba(209, 39, 48, 0.04);\n border-color: rgba(209, 39, 48, 0.06);\n}\n\na.tb-item-card.tb-active {\n background: rgba(48, 86, 128, 0.04);\n border-color: rgba(48, 86, 128, 0.12);\n}\n\na.tb-item-card.tb-total {\n background: rgba(0, 0, 0, 0.01);\n border-color: rgba(0, 0, 0, 0.05);\n}\n\n.tb-item-title-container {\n display: grid;\n}\n\n.tb-item-title {\n font-weight: 400;\n font-size: 14px;\n line-height: 20px;\n letter-spacing: 0.2px;\n white-space: nowrap;\n overflow: hidden;\n text-overflow: ellipsis; \n color: rgba(0, 0, 0, 0.76);\n}\n\n.tb-item-title.tb-home-widget-link:after {\n position: absolute;\n right: 0;\n}\n\na.tb-item-card:hover .tb-item-title.tb-home-widget-link:after { \n color: rgba(0, 0, 0, 0.38);\n}\n\na.tb-item-card:hover {\n box-shadow: 0px 4px 10px rgba(23, 33, 90, 0.08);\n}\n\n.tb-count-container {\n flex: 1;\n display: flex;\n flex-direction: column;\n align-items: flex-start;\n justify-content: center;\n}\n\n.tb-count {\n font-style: normal;\n font-weight: 500;\n font-size: 24px;\n line-height: 36px;\n white-space: nowrap;\n color: rgba(0, 0, 0, 0.87);\n}\n\n@media screen and (max-width: 959px) {\n .tb-item-cards {\n flex-direction: column;\n }\n a.tb-item-card {\n margin-bottom: 0;\n }\n}\n\n@media screen and (max-width: 1279px) {\n a.tb-item-card {\n flex-direction: row;\n align-items: center;\n }\n .tb-item-title.tb-home-widget-link:after {\n position: relative;\n }\n .tb-count-container {\n align-items: flex-end;\n }\n}\n\n@media screen and (min-width: 960px) and (max-width: 1819px) {\n .tb-item-title {\n font-size: 11px;\n line-height: 16px;\n }\n .tb-count {\n font-size: 16px;\n line-height: 24px;\n }\n a.tb-item-card {\n padding: 4px 8px;\n margin-bottom: 6px;\n }\n a.tb-item-card:hover {\n box-shadow: 0px 2px 5px rgba(23, 33, 90, 0.08);\n }\n}\n"
},

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

@ -1102,6 +1102,7 @@
"expression": "Expression",
"no-found": "No calculated fields found",
"list": "{ count, plural, =1 {One calculated field} other {List of # calculated fields} }",
"add": "Add calculated field",
"selected-fields": "{ count, plural, =1 {1 calculated field} other {# calculated fields} } selected",
"type": {
"simple": "Simple",

Loading…
Cancel
Save