From 6cf314eaa01ce62d9ffbaba02e05d77a722817fc Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 19 Feb 2025 08:35:24 +0200 Subject: [PATCH 01/13] added NaN value when is not number in rolling --- .../server/service/cf/ctx/state/TsRollingArgumentEntry.java | 6 +++--- .../org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java | 3 +++ .../server/dao/service/CalculatedFieldServiceTest.java | 6 +++--- .../service/validator/CalculatedFieldDataValidatorTest.java | 3 +++ 4 files changed, 12 insertions(+), 6 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index 866e2f5e09..da02b9be2e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -86,7 +86,7 @@ public class TsRollingArgumentEntry implements ArgumentEntry { } @Override - public boolean updateEntry(ArgumentEntry entry) throws CalculatedFieldStateException { + public boolean updateEntry(ArgumentEntry entry) { if (entry instanceof TsRollingArgumentEntry tsRollingEntry) { updateTsRollingEntry(tsRollingEntry); } else if (entry instanceof SingleValueArgumentEntry singleValueEntry) { @@ -118,8 +118,8 @@ public class TsRollingArgumentEntry implements ArgumentEntry { } cleanupExpiredRecords(); } catch (Exception e) { - log.warn("Time series rolling arguments supports only numeric values."); -// throw new IllegalArgumentException("Time series rolling arguments supports only numeric values."); + tsRecords.put(ts, Double.NaN); + log.warn("Invalid value '{}' for time series rolling arguments. Only numeric values are supported.", value.getValue()); } } diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java index 843abf8d5c..c3a3606ad1 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java @@ -63,6 +63,9 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable Date: Wed, 19 Feb 2025 08:59:42 +0200 Subject: [PATCH 02/13] fixed tests --- .../server/queue/discovery/HashPartitionServiceTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java index 52cbaa4add..3a015d23b2 100644 --- a/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java @@ -426,6 +426,9 @@ public class HashPartitionServiceTest { topicService); ReflectionTestUtils.setField(partitionService, "coreTopic", "tb.core"); ReflectionTestUtils.setField(partitionService, "corePartitions", 10); + ReflectionTestUtils.setField(partitionService, "cfEventTopic", "tb_cf_event"); + ReflectionTestUtils.setField(partitionService, "cfStateTopic", "tb_cf_state"); + ReflectionTestUtils.setField(partitionService, "cfPartitions", 10); ReflectionTestUtils.setField(partitionService, "vcTopic", "tb.vc"); ReflectionTestUtils.setField(partitionService, "vcPartitions", 10); ReflectionTestUtils.setField(partitionService, "hashFunctionName", hashFunctionName); From 300650ceb6eaa2811db129be2069a23d59af721c Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 19 Feb 2025 09:06:23 +0200 Subject: [PATCH 03/13] fixed test where rolling arg is NaN --- .../service/cf/ctx/state/ScriptCalculatedFieldStateTest.java | 2 +- .../service/cf/ctx/state/TsRollingArgumentEntryTest.java | 5 ++--- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java index bc75bd74ff..e52e6e985a 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java @@ -151,7 +151,7 @@ public class ScriptCalculatedFieldStateTest { } private TsRollingArgumentEntry createRollingArgEntry() { - TsRollingArgumentEntry argumentEntry = new TsRollingArgumentEntry(); + TsRollingArgumentEntry argumentEntry = new TsRollingArgumentEntry(5, 30000L); long ts = System.currentTimeMillis(); TreeMap values = new TreeMap<>(); diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java index 75aca61825..8f8abadf02 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntryTest.java @@ -79,9 +79,8 @@ public class TsRollingArgumentEntryTest { void testUpdateEntryWhenValueIsNotNumber() { SingleValueArgumentEntry newEntry = new SingleValueArgumentEntry(ts - 10, new StringDataEntry("key", "string"), 123L); - assertThatThrownBy(() -> entry.updateEntry(newEntry)) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Time series rolling arguments supports only numeric values."); + assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.getTsRecords().get(ts - 10)).isNaN(); } @Test From 609a5dc34b67694620d8cc90620f85acf46ca9f8 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 19 Feb 2025 10:18:41 +0200 Subject: [PATCH 04/13] fixed cache onCfDelete and updated tenant profile properties --- ...CalculatedFieldEntityMessageProcessor.java | 29 ++++++++++--------- ...alculatedFieldManagerMessageProcessor.java | 1 + .../DefaultTenantProfileConfiguration.java | 3 +- .../CalculatedFieldDataValidator.java | 14 +++++++-- 4 files changed, 30 insertions(+), 17 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 7211ac0db9..7b7ce63d38 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -187,23 +187,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } } catch (Exception e) { + if (e instanceof CalculatedFieldException) { + throw (CalculatedFieldException) e; + } throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); } } - @SneakyThrows - private void processTelemetry(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List cfIdList, MultipleTbCallback callback) { + private void processTelemetry(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List cfIdList, MultipleTbCallback callback) throws CalculatedFieldException { processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); } - @SneakyThrows - private void processAttributes(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List cfIdList, MultipleTbCallback callback) { + private void processAttributes(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List cfIdList, MultipleTbCallback callback) throws CalculatedFieldException { processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto)); } - @SneakyThrows private void processArgumentValuesUpdate(CalculatedFieldCtx ctx, List cfIdList, MultipleTbCallback callback, - Map newArgValues, UUID tbMsgId, TbMsgType tbMsgType) { + Map newArgValues, UUID tbMsgId, TbMsgType tbMsgType) throws CalculatedFieldException { if (newArgValues.isEmpty()) { log.info("[{}] No new argument values to process for CF.", ctx.getCfId()); callback.onSuccess(CALLBACKS_PER_CF); @@ -241,15 +241,18 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM return state; } - @SneakyThrows - private void processStateIfReady(CalculatedFieldCtx ctx, List cfIdList, CalculatedFieldState state, UUID tbMsgId, TbMsgType tbMsgType, TbCallback callback) { + private void processStateIfReady(CalculatedFieldCtx ctx, List cfIdList, CalculatedFieldState state, UUID tbMsgId, TbMsgType tbMsgType, TbCallback callback) throws CalculatedFieldException { CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId); if (state.isReady() && ctx.isInitialized()) { - CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(5, TimeUnit.SECONDS); - state.checkStateSize(ctxId, ctx.getMaxStateSizeInKBytes()); - cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback); - if (DebugModeUtil.isDebugAllAvailable(ctx.getCalculatedField())) { - systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, JacksonUtil.writeValueAsString(calculationResult.getResult()), null); + try { + CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(5, TimeUnit.SECONDS); + state.checkStateSize(ctxId, ctx.getMaxStateSizeInKBytes()); + cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback); + if (DebugModeUtil.isDebugAllAvailable(ctx.getCalculatedField())) { + systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, JacksonUtil.writeValueAsString(calculationResult.getResult()), null); + } + } catch (Exception e) { + throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).msgId(tbMsgId).msgType(tbMsgType).arguments(state.getArguments()).cause(e).build(); } } else { state.checkStateSize(ctxId, ctx.getMaxStateSizeInKBytes()); diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 5e005935aa..1d77f7cd76 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -308,6 +308,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware log.warn("[{}] CF was already deleted [{}]", tenantId, cfId); callback.onSuccess(); } else { + entityIdCalculatedFields.get(cfCtx.getEntityId()).remove(cfCtx); deleteLinks(cfCtx); EntityId entityId = cfCtx.getEntityId(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 0139792fa7..0b7b030748 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -135,7 +135,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private double warnThreshold; - private long maxCalculatedFields; + private long maxCalculatedFieldsPerEntity; private long maxArgumentsPerCF; private long maxDataPointsPerRollingArg; private long maxStateSizeInKBytes; @@ -182,7 +182,6 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura case DASHBOARD -> maxDashboards; case RULE_CHAIN -> maxRuleChains; case EDGE -> maxEdges; - case CALCULATED_FIELD -> maxCalculatedFields; default -> 0; }; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java index b52c2fe9b7..12c7e15f49 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java @@ -17,8 +17,8 @@ package org.thingsboard.server.dao.service.validator; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.dao.cf.CalculatedFieldDao; @@ -37,7 +37,7 @@ public class CalculatedFieldDataValidator extends DataValidator @Override protected void validateCreate(TenantId tenantId, CalculatedField calculatedField) { - validateNumberOfEntitiesPerTenant(tenantId, EntityType.CALCULATED_FIELD); + validateNumberOfCFsPerEntity(tenantId, calculatedField.getEntityId()); validateNumberOfArgumentsPerCF(tenantId, calculatedField); } @@ -51,6 +51,16 @@ public class CalculatedFieldDataValidator extends DataValidator return old; } + private void validateNumberOfCFsPerEntity(TenantId tenantId, EntityId entityId) { + long maxCFsPerEntity = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxCalculatedFieldsPerEntity); + if (maxCFsPerEntity <= 0) { + return; + } + if (calculatedFieldDao.countCFByEntityId(tenantId, entityId) >= maxCFsPerEntity) { + throw new DataValidationException("Calculated fields per entity limit reached!"); + } + } + private void validateNumberOfArgumentsPerCF(TenantId tenantId, CalculatedField calculatedField) { long maxArgumentsPerCF = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxArgumentsPerCF); if (maxArgumentsPerCF <= 0) { From e018282d0410349531a92422a23025d839b588bb Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 19 Feb 2025 17:51:48 +0200 Subject: [PATCH 05/13] minor fixes --- .../CalculatedFieldEntityMessageProcessor.java | 3 +++ .../CalculatedFieldManagerMessageProcessor.java | 13 +++---------- .../controller/CalculatedFieldControllerTest.java | 5 +++-- .../org/thingsboard/script/api/tbel/TbelCfArg.java | 2 ++ 4 files changed, 11 insertions(+), 12 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 7b7ce63d38..e551d9f067 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -116,6 +116,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM var cfState = getOrInitState(cfCtx); processStateIfReady(cfCtx, Collections.singletonList(cfCtx.getCfId()), cfState, null, null, msg.getCallback()); } catch (Exception e) { + if (e instanceof CalculatedFieldException cfe) { + throw cfe; + } throw CalculatedFieldException.builder().ctx(cfCtx).eventEntity(entityId).cause(e).build(); } } diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 1d77f7cd76..b65a4072d3 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldInitMsg; @@ -39,7 +38,6 @@ import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.dao.cf.CalculatedFieldService; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityCtxIdProto; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldStateService; import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache; @@ -53,11 +51,12 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CopyOnWriteArrayList; +import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; + /** * @author Andrew Shvayka @@ -358,7 +357,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware var proto = msg.getProto(); var linksList = proto.getLinksList(); for (var linkProto : linksList) { - var link = toCalculatedFieldEntityCtxId(linkProto); + var link = fromProto(linkProto); var targetEntityId = link.entityId(); var targetEntityType = targetEntityId.getEntityType(); var cf = calculatedFields.get(link.cfId()); @@ -383,12 +382,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } } - private CalculatedFieldEntityCtxId toCalculatedFieldEntityCtxId(CalculatedFieldEntityCtxIdProto ctxIdProto) { - EntityId entityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); - CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB())); - return new CalculatedFieldEntityCtxId(tenantId, calculatedFieldId, entityId); - } - private List filterCalculatedFieldLinks(CalculatedFieldTelemetryMsg msg) { EntityId entityId = msg.getEntityId(); var proto = msg.getProto(); diff --git a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java index b1c7547251..79288fdb95 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -86,13 +86,14 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { assertThat(savedCalculatedField.getType()).isEqualTo(calculatedField.getType()); assertThat(savedCalculatedField.getName()).isEqualTo(calculatedField.getName()); assertThat(savedCalculatedField.getConfiguration()).isEqualTo(getCalculatedFieldConfig(testDevice.getId())); - assertThat(savedCalculatedField.getVersion()).isEqualTo(calculatedField.getVersion()); + assertThat(savedCalculatedField.getVersion()).isEqualTo(1L); savedCalculatedField.setName("Test CF"); CalculatedField updatedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); - assertThat(updatedCalculatedField).isEqualTo(savedCalculatedField); + assertThat(updatedCalculatedField.getName()).isEqualTo(savedCalculatedField.getName()); + assertThat(updatedCalculatedField.getVersion()).isEqualTo(savedCalculatedField.getVersion() + 1); doDelete("/api/calculatedField/" + savedCalculatedField.getId().getId().toString()) .andExpect(status().isOk()); diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfArg.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfArg.java index 7fe75acd05..ddeb9d14c4 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfArg.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfArg.java @@ -15,6 +15,7 @@ */ package org.thingsboard.script.api.tbel; +import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; @@ -29,6 +30,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; }) public interface TbelCfArg extends TbelCfObject { + @JsonIgnore String getType(); } From 40491f009e52d889ad6d501a34885e50e909a480 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 09:43:24 +0200 Subject: [PATCH 06/13] removed cf debug events rate limit from tenant profile --- .../cf/CalculatedFieldIntegrationTest.java | 54 +++++++------------ .../server/common/data/limit/LimitedApi.java | 2 +- .../DefaultTenantProfileConfiguration.java | 1 - 3 files changed, 19 insertions(+), 38 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java index f544f463fa..6e90d51d35 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -83,57 +83,52 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes calculatedField.setConfiguration(config); calculatedField.setVersion(1L); - // create CF -> perform initial calculation CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("77.0"); }); - // update telemetry -> recalculate state doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"temperature\":30}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update telemetry -> recalculate state").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); }); - // update CF output -> perform calculation with updated output Output savedOutput = savedCalculatedField.getConfiguration().getOutput(); savedOutput.setType(OutputType.ATTRIBUTES); savedOutput.setScope(AttributeScope.SERVER_SCOPE); savedOutput.setName("temperatureF"); savedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update CF output -> perform calculation with updated output").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ArrayNode temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); assertThat(temperatureF).isNotNull(); assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("86.0"); }); - // update CF argument -> perform calculation with new argument Argument savedArgument = savedCalculatedField.getConfiguration().getArguments().get("T"); savedArgument.setRefEntityKey(new ReferencedEntityKey("deviceTemperature", ArgumentType.ATTRIBUTE, AttributeScope.SERVER_SCOPE)); savedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update CF argument -> perform calculation with new argument").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ArrayNode temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); assertThat(temperatureF).isNotNull(); assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("104.0"); }); - // update CF expression -> perform calculation with new expression savedCalculatedField.getConfiguration().setExpression("1.8 * T + 32"); savedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update CF expression -> perform calculation with new expression").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ArrayNode temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); assertThat(temperatureF).isNotNull(); @@ -168,20 +163,18 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes calculatedField.setConfiguration(config); calculatedField.setVersion(1L); - // create CF -> state is not ready -> no calculation performed CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("create CF -> state is not ready -> no calculation performed").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); }); - // update telemetry -> perform calculation doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"temperature\":30}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update telemetry -> perform calculation").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); @@ -217,20 +210,18 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes calculatedField.setConfiguration(config); calculatedField.setVersion(1L); - // create CF -> perform initial calculation with default value CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("create CF -> perform initial calculation with default value").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("53.6"); }); - // update telemetry -> recalculate state doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"temperature\":30}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update telemetry -> recalculate state").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); @@ -283,10 +274,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes calculatedField.setConfiguration(config); calculatedField.setVersion(1L); - // create CF and perform initial calculation doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("create CF and perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -299,10 +289,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(z2.get(0).get("value").asText()).isEqualTo("52.0"); }); - // update device telemetry -> recalculate state for all assets doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"x\":25}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update device telemetry -> recalculate state for all assets").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -315,10 +304,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(z2.get(0).get("value").asText()).isEqualTo("37.0"); }); - // update asset 1 telemetry -> recalculate state only for asset 1 doPost("/api/plugins/telemetry/ASSET/" + asset1.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"y\":15}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update asset 1 telemetry -> recalculate state only for asset 1").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -331,10 +319,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(z2.get(0).get("value").asText()).isEqualTo("37.0"); }); - // update asset 2 telemetry -> recalculate state only for asset 2 doPost("/api/plugins/telemetry/ASSET/" + asset2.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"y\":5}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update asset 2 telemetry -> recalculate state only for asset 2").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 (no changes) ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -347,12 +334,11 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(z2.get(0).get("value").asText()).isEqualTo("30.0"); }); - // add new entity to profile -> calculate state for new entity Asset asset3 = createAsset("Test asset 3", assetProfile.getId()); doPost("/api/plugins/telemetry/ASSET/" + asset3.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"y\":13}")); Asset finalAsset3 = asset3; - await().atMost(5, TimeUnit.SECONDS) + await().alias("add new entity to profile -> calculate state for new entity").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 3 ArrayNode z3 = getServerAttributes(finalAsset3.getId(), "z"); @@ -360,10 +346,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(z3.get(0).get("value").asText()).isEqualTo("38.0"); }); - // update device telemetry -> recalculate state for all assets doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"x\":20}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update device telemetry -> recalculate state for all assets").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -386,11 +371,10 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes asset3.setAssetProfileId(newAssetProfile.getId()); asset3 = doPost("/api/asset", asset3, Asset.class); - // update device telemetry -> recalculate state for asset 1 and asset 2 doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/attributes/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"x\":15}")); Asset updatedAsset3 = asset3; - await().atMost(5, TimeUnit.SECONDS) + await().alias("update device telemetry -> recalculate state for asset 1 and asset 2").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { // result of asset 1 ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); @@ -438,20 +422,18 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes calculatedField.setConfiguration(config); calculatedField.setVersion(1L); - // create CF -> ctx is not initialized -> no calculation perform CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().atMost(5, TimeUnit.SECONDS) + await().alias("create CF -> ctx is not initialized -> no calculation perform").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); }); - // update telemetry -> ctx is not initialized -> no calculation perform doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"temperature\":30}")); - await().atMost(5, TimeUnit.SECONDS) + await().alias("update telemetry -> ctx is not initialized -> no calculation perform").atMost(TIMEOUT, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java index 35068821ee..141a805dc0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java @@ -44,7 +44,7 @@ public enum LimitedApi { TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE("transport messages per gateway device", false), EMAILS("emails sending", true), WS_SUBSCRIPTIONS("WS subscriptions", false), - CALCULATED_FIELD_DEBUG_EVENTS(DefaultTenantProfileConfiguration::getCalculatedFieldDebugEventsRateLimit, "calculated field debug events", true); + CALCULATED_FIELD_DEBUG_EVENTS("calculated field debug events", true); private Function configExtractor; @Getter diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 0b7b030748..464ab410ab 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -140,7 +140,6 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxDataPointsPerRollingArg; private long maxStateSizeInKBytes; private long maxSingleValueArgumentSizeInKBytes; - private String calculatedFieldDebugEventsRateLimit; @Override public long getProfileThreshold(ApiUsageRecordKey key) { From 6f1dd5a2d6d4041925e59c13df007afe9f042b24 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 10:58:26 +0200 Subject: [PATCH 07/13] added new methods for working with rolling --- .../script/api/tbel/TbelCfTsRollingArg.java | 119 ++++++++++++++++++ 1 file changed, 119 insertions(+) diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java index c3a3606ad1..ceae4e60c5 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java @@ -20,6 +20,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import lombok.Getter; +import java.util.ArrayList; import java.util.Collections; import java.util.Iterator; import java.util.List; @@ -60,6 +61,10 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable val) { + min = val; + } + } + return min; + } + + public double mean() { + if (values.isEmpty()) { + return 0; + } + + double sum = sum(); + return Double.isNaN(sum) ? Double.NaN : sum / values.size(); + } + + public double std() { + if (values.isEmpty()) { + return 0; + } + + double mean = mean(); + if (Double.isNaN(mean)) { + return Double.NaN; + } + + double sum = 0; + for (TbelCfTsDoubleVal value : values) { + double val = value.getValue(); + if (Double.isNaN(val)) { + return Double.NaN; + } + sum += Math.pow(val - mean, 2); + } + return Math.sqrt(sum / values.size()); + } + + public double median() { + if (values.isEmpty()) { + return 0; + } + + List sortedValues = new ArrayList<>(); + for (TbelCfTsDoubleVal value : values) { + double val = value.getValue(); + if (Double.isNaN(val)) { + return Double.NaN; + } + sortedValues.add(val); + } + Collections.sort(sortedValues); + + int size = sortedValues.size(); + return (size % 2 == 1) + ? sortedValues.get(size / 2) + : (sortedValues.get(size / 2 - 1) + sortedValues.get(size / 2)) / 2.0; + } + + public double count() { + return hasNaN() ? Double.NaN : values.size(); + } + + public double last() { + if (values.isEmpty()) { + return 0; + } + + return hasNaN() ? Double.NaN : values.get(values.size() - 1).getValue(); + } + + public double first() { + if (values.isEmpty()) { + return 0; + } + + return hasNaN() ? Double.NaN : values.get(0).getValue(); + } + + public double sum() { + if (values.isEmpty()) { + return 0; + } + + double sum = 0; + for (TbelCfTsDoubleVal value : values) { + double val = value.getValue(); + if (Double.isNaN(val)) { + return Double.NaN; + } + sum += val; + } + return sum; + } + + private boolean hasNaN() { + for (TbelCfTsDoubleVal value : values) { + if (Double.isNaN(value.getValue())) { + return true; + } + } + return false; + } + @JsonIgnore public int getSize() { return values.size(); From 6698a8f8f5ae530de6299f9c76b9ba5ddfca5cc1 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 11:33:56 +0200 Subject: [PATCH 08/13] fixed tenant actor test --- .../org/thingsboard/server/actors/tenant/TenantActorTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/actors/tenant/TenantActorTest.java b/application/src/test/java/org/thingsboard/server/actors/tenant/TenantActorTest.java index cad933ccc5..ed8468aa6a 100644 --- a/application/src/test/java/org/thingsboard/server/actors/tenant/TenantActorTest.java +++ b/application/src/test/java/org/thingsboard/server/actors/tenant/TenantActorTest.java @@ -27,6 +27,7 @@ import org.thingsboard.server.actors.TbActorSystemSettings; import org.thingsboard.server.actors.TbEntityActorId; import org.thingsboard.server.actors.ruleChain.RuleChainActor; import org.thingsboard.server.actors.ruleChain.RuleChainToRuleChainMsg; +import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.shared.RuleChainErrorActor; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.id.DeviceId; @@ -116,6 +117,7 @@ public class TenantActorTest { TbActorSystemSettings settings = new TbActorSystemSettings(0, 0, 0); TbActorSystem system = spy(new DefaultTbActorSystem(settings)); system.createDispatcher(RULE_DISPATCHER_NAME, mock()); + system.createDispatcher(DefaultActorService.CF_MANAGER_DISPATCHER_NAME, mock()); TbActorMailbox tenantCtx = new TbActorMailbox(system, settings, null, mock(), mock(), null); tenantActor.init(tenantCtx); From a6a29513f618a826aa7c62cb0f46ca6954e24f44 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 11:50:34 +0200 Subject: [PATCH 09/13] removed broadcast asset updates event to transport --- .../server/service/queue/DefaultTbClusterService.java | 1 - 1 file changed, 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index 920a7563dc..19c386fff3 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -697,7 +697,6 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void onAssetUpdated(Asset entity, Asset old) { var created = old == null; - broadcastEntityChangeToTransport(entity.getTenantId(), entity.getId(), entity, null); if (old != null) { boolean assetTypeChanged = !entity.getAssetProfileId().equals(old.getAssetProfileId()); if (assetTypeChanged) { From f362749a7b24c5f9c32fab7a38ea0cb97db8fc0a Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 12:09:25 +0200 Subject: [PATCH 10/13] fixed tenant profile controller test --- .../server/controller/TenantProfileControllerTest.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java index 989e80e6e5..ad8eb188c9 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/TenantProfileControllerTest.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfiguration; import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.queue.TbQueueCallback; import java.util.ArrayList; import java.util.Collections; @@ -44,6 +45,7 @@ import java.util.List; import java.util.stream.Collectors; import static org.hamcrest.Matchers.containsString; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -354,7 +356,7 @@ public class TenantProfileControllerTest extends AbstractControllerTest { argument -> argument.getClass().equals(TenantProfile.class); if (ComponentLifecycleEvent.DELETED.equals(event)) { Mockito.verify(tbClusterService, times(cntTime)).onTenantProfileDelete(Mockito.argThat(matcherTenantProfile), - Mockito.isNull()); + eq(TbQueueCallback.EMPTY)); testBroadcastEntityStateChangeEventNever(createEntityId_NULL_UUID(new Tenant())); } else { Mockito.verify(tbClusterService, times(cntTime)).onTenantProfileChange(Mockito.argThat(matcherTenantProfile), From a52c578ecdad303a8bdd2b8a170d77e69facb2fd Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 25 Feb 2025 15:06:40 +0200 Subject: [PATCH 11/13] minor fixes --- .../CalculatedFieldEntityMessageProcessor.java | 4 ++-- .../server/service/cf/ctx/state/TsRollingArgumentEntry.java | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index ecf1b140f2..e564489da2 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -194,8 +194,8 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } } catch (Exception e) { - if (e instanceof CalculatedFieldException) { - throw (CalculatedFieldException) e; + if (e instanceof CalculatedFieldException cfe) { + throw cfe; } throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index da02b9be2e..90068d6b6f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -119,7 +119,7 @@ public class TsRollingArgumentEntry implements ArgumentEntry { cleanupExpiredRecords(); } catch (Exception e) { tsRecords.put(ts, Double.NaN); - log.warn("Invalid value '{}' for time series rolling arguments. Only numeric values are supported.", value.getValue()); + log.debug("Invalid value '{}' for time series rolling arguments. Only numeric values are supported.", value.getValue()); } } From 72f26d9fe2b22bbd14fe94ca702c0bc462d2d26f Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 26 Feb 2025 15:22:42 +0200 Subject: [PATCH 12/13] fixed functions --- .../cf/ctx/state/TsRollingArgumentEntry.java | 4 +- .../script/api/tbel/TbelCfTsRollingArg.java | 132 +++++++++++++----- 2 files changed, 100 insertions(+), 36 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index 90068d6b6f..bb559a3eec 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -25,7 +25,6 @@ import org.thingsboard.script.api.tbel.TbelCfTsDoubleVal; import org.thingsboard.script.api.tbel.TbelCfTsRollingArg; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; -import org.thingsboard.server.exception.CalculatedFieldStateException; import java.util.ArrayList; import java.util.List; @@ -116,10 +115,11 @@ public class TsRollingArgumentEntry implements ArgumentEntry { case STRING -> value.getStrValue().ifPresent(aString -> tsRecords.put(ts, Double.parseDouble(aString))); case JSON -> value.getJsonValue().ifPresent(aString -> tsRecords.put(ts, Double.parseDouble(aString))); } - cleanupExpiredRecords(); } catch (Exception e) { tsRecords.put(ts, Double.NaN); log.debug("Invalid value '{}' for time series rolling arguments. Only numeric values are supported.", value.getValue()); + } finally { + cleanupExpiredRecords(); } } diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java index ceae4e60c5..1bae5399cc 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArg.java @@ -61,14 +61,18 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable val) { @@ -97,21 +105,28 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable sortedValues = new ArrayList<>(); for (TbelCfTsDoubleVal value : values) { double val = value.getValue(); if (Double.isNaN(val)) { - return Double.NaN; + if (!ignoreNaN) { + return Double.NaN; + } + } else { + sortedValues.add(val); } - sortedValues.add(val); } Collections.sort(sortedValues); @@ -147,51 +172,90 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable= 0; i--) { + double prevValue = values.get(i).getValue(); + if (!Double.isNaN(prevValue)) { + return prevValue; + } + } + throw new IllegalArgumentException("Rolling argument values are empty."); } public double first() { + return first(true); + } + + public double first(boolean ignoreNaN) { if (values.isEmpty()) { - return 0; + throw new IllegalArgumentException("Rolling argument values are empty."); } - return hasNaN() ? Double.NaN : values.get(0).getValue(); + double firstValue = values.get(0).getValue(); + if (!Double.isNaN(firstValue) || !ignoreNaN) { + return firstValue; + } + for (int i = 1; i < values.size(); i++) { + double nextValue = values.get(i).getValue(); + if (!Double.isNaN(nextValue)) { + return nextValue; + } + } + throw new IllegalArgumentException("Rolling argument values are empty."); } public double sum() { + return sum(true); + } + + public double sum(boolean ignoreNaN) { if (values.isEmpty()) { - return 0; + throw new IllegalArgumentException("Rolling argument values are empty."); } double sum = 0; for (TbelCfTsDoubleVal value : values) { double val = value.getValue(); if (Double.isNaN(val)) { - return Double.NaN; + if (!ignoreNaN) { + return Double.NaN; + } + } else { + sum += val; } - sum += val; } return sum; } - private boolean hasNaN() { - for (TbelCfTsDoubleVal value : values) { - if (Double.isNaN(value.getValue())) { - return true; - } - } - return false; - } - @JsonIgnore public int getSize() { return values.size(); From 8c86f7fc12cf686ce61e47063875b60ca4a18c90 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 26 Feb 2025 15:48:12 +0200 Subject: [PATCH 13/13] added tests for rolling methods --- .../api/tbel/TbelCfTsRollingArgTest.java | 131 ++++++++++++++++++ 1 file changed, 131 insertions(+) create mode 100644 common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArgTest.java diff --git a/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArgTest.java b/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArgTest.java new file mode 100644 index 0000000000..e4ebf8a8c1 --- /dev/null +++ b/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbelCfTsRollingArgTest.java @@ -0,0 +1,131 @@ +/** + * Copyright © 2016-2024 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.script.api.tbel; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.assertj.core.api.Assertions.within; + +public class TbelCfTsRollingArgTest { + + private final long ts = System.currentTimeMillis(); + + private TbelCfTsRollingArg rollingArg; + + @BeforeEach + void setUp() { + rollingArg = new TbelCfTsRollingArg( + new TbTimeWindow(ts - 30000, ts - 10, 10), + List.of( + new TbelCfTsDoubleVal(ts - 10, Double.NaN), + new TbelCfTsDoubleVal(ts - 20, 2.0), + new TbelCfTsDoubleVal(ts - 30, 8.0), + new TbelCfTsDoubleVal(ts - 40, Double.NaN), + new TbelCfTsDoubleVal(ts - 50, 3.0), + new TbelCfTsDoubleVal(ts - 60, 9.0), + new TbelCfTsDoubleVal(ts - 70, Double.NaN) + ) + ); + } + + @Test + void testMax() { + assertThat(rollingArg.max()).isEqualTo(9.0); + assertThat(rollingArg.max(false)).isNaN(); + } + + @Test + void testMin() { + assertThat(rollingArg.min()).isEqualTo(2.0); + assertThat(rollingArg.min(false)).isNaN(); + } + + @Test + void testMean() { + assertThat(rollingArg.mean()).isEqualTo(5.5); + assertThat(rollingArg.mean(false)).isNaN(); + } + + @Test + void testStd() { + assertThat(rollingArg.std()).isCloseTo(3.0413812651491097, within(0.001)); + assertThat(rollingArg.std(false)).isNaN(); + } + + @Test + void testMedian() { + assertThat(rollingArg.median()).isEqualTo(5.5); + assertThat(rollingArg.median(false)).isNaN(); + } + + @Test + void testCount() { + assertThat(rollingArg.count()).isEqualTo(4); + assertThat(rollingArg.count(false)).isEqualTo(7); + } + + @Test + void testLast() { + assertThat(rollingArg.last()).isEqualTo(9.0); + assertThat(rollingArg.last(false)).isNaN(); + } + + @Test + void testFirst() { + assertThat(rollingArg.first()).isEqualTo(2.0); + assertThat(rollingArg.first(false)).isNaN(); + } + + @Test + void testFirstAndLastWhenOnlyNaNAndIgnoreNaNIsFalse() { + assertThat(rollingArg.first()).isEqualTo(2.0); + rollingArg = new TbelCfTsRollingArg( + new TbTimeWindow(ts - 30000, ts - 10, 10), + List.of( + new TbelCfTsDoubleVal(ts - 10, Double.NaN), + new TbelCfTsDoubleVal(ts - 40, Double.NaN), + new TbelCfTsDoubleVal(ts - 70, Double.NaN) + ) + ); + assertThatThrownBy(rollingArg::first).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::last).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + } + + @Test + void testSum() { + assertThat(rollingArg.sum()).isEqualTo(22.0); + assertThat(rollingArg.sum(false)).isNaN(); + } + + @Test + void testEmptyValues() { + rollingArg = new TbelCfTsRollingArg(new TbTimeWindow(0, 10, 10), List.of()); + assertThatThrownBy(rollingArg::sum).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::max).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::min).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::mean).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::std).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::median).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::first).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + assertThatThrownBy(rollingArg::last).isInstanceOf(IllegalArgumentException.class).hasMessage("Rolling argument values are empty."); + } + +} \ No newline at end of file