diff --git a/.github/workflows/check-configuration-files.yml b/.github/workflows/check-configuration-files.yml
index f280e326c2..561b7d0019 100644
--- a/.github/workflows/check-configuration-files.yml
+++ b/.github/workflows/check-configuration-files.yml
@@ -32,14 +32,14 @@ on:
jobs:
build:
name: Check thingsboard.yml file
- runs-on: ubuntu-20.04
+ runs-on: ubuntu-22.04
steps:
- name: Checkout code
uses: actions/checkout@v2
- - name: Set up Python 3.10
+ - name: Set up Python 3.13
uses: actions/setup-python@v3
with:
- python-version: "3.10.2"
+ python-version: "3.13.2"
architecture: "x64"
env:
AGENT_TOOLSDIRECTORY: /opt/hostedtoolcache
diff --git a/application/pom.xml b/application/pom.xml
index c6058408f8..6ce77646f8 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -124,6 +124,10 @@
org.thingsboard.common
edge-api
+
+ org.thingsboard.common
+ edqs
+
org.thingsboard
dao
diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql
index ddfac0d768..e9cbf8ef8e 100644
--- a/application/src/main/data/upgrade/basic/schema_update.sql
+++ b/application/src/main/data/upgrade/basic/schema_update.sql
@@ -63,4 +63,21 @@ $$;
-- UPDATE SAVE TIME SERIES NODES END
-ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS version BIGINT DEFAULT 1;
\ No newline at end of file
+ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS version BIGINT DEFAULT 1;
+
+-- UPDATE TENANT PROFILE CALCULATED FIELD LIMITS START
+
+UPDATE tenant_profile
+SET profile_data = profile_data
+ || jsonb_build_object(
+ 'configuration', profile_data->'configuration' || jsonb_build_object(
+ 'maxCalculatedFieldsPerEntity', COALESCE(profile_data->'configuration'->>'maxCalculatedFieldsPerEntity', '5')::bigint,
+ 'maxArgumentsPerCF', COALESCE(profile_data->'configuration'->>'maxArgumentsPerCF', '10')::bigint,
+ 'maxDataPointsPerRollingArg', COALESCE(profile_data->'configuration'->>'maxDataPointsPerRollingArg', '1000')::bigint,
+ 'maxStateSizeInKBytes', COALESCE(profile_data->'configuration'->>'maxStateSizeInKBytes', '32')::bigint,
+ 'maxSingleValueArgumentSizeInKBytes', COALESCE(profile_data->'configuration'->>'maxSingleValueArgumentSizeInKBytes', '2')::bigint
+ )
+ )
+WHERE profile_data->'configuration'->>'maxCalculatedFieldsPerEntity' IS NULL;
+
+-- UPDATE TENANT PROFILE CALCULATED FIELD LIMITS END
diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
index 50a777840c..2c22941c2f 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
@@ -33,7 +33,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.MailService;
import org.thingsboard.rule.engine.api.NotificationCenter;
-import org.thingsboard.rule.engine.api.RuleEngineDeviceStateManager;
+import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.SmsService;
import org.thingsboard.rule.engine.api.notification.SlackService;
import org.thingsboard.rule.engine.api.sms.SmsSenderFactory;
@@ -109,6 +109,7 @@ import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.cf.CalculatedFieldProcessingService;
+import org.thingsboard.server.service.cf.CalculatedFieldQueueService;
import org.thingsboard.server.service.cf.CalculatedFieldStateService;
import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
@@ -233,7 +234,7 @@ public class ActorSystemContext {
@Autowired(required = false)
@Getter
- private RuleEngineDeviceStateManager deviceStateManager;
+ private DeviceStateManager deviceStateManager;
@Autowired
@Getter
@@ -544,6 +545,11 @@ public class ActorSystemContext {
@Getter
private CalculatedFieldStateService calculatedFieldStateService;
+ @Lazy
+ @Autowired(required = false)
+ @Getter
+ private CalculatedFieldQueueService calculatedFieldQueueService;
+
@Lazy
@Autowired(required = false)
@Getter
@@ -644,6 +650,10 @@ public class ActorSystemContext {
@Getter
private String deviceStateNodeRateLimitConfig;
+ @Value("${actors.calculated_fields.calculation_timeout:5}")
+ @Getter
+ private long cfCalculationResultTimeout;
+
@Getter
@Setter
private TbActorSystem actorSystem;
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 f7fc204c0f..9ab019097b 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
@@ -134,14 +134,20 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
public void process(CalculatedFieldEntityDeleteMsg msg) {
log.info("[{}] Processing CF entity delete msg.", msg.getEntityId());
if (this.entityId.equals(msg.getEntityId())) {
- MultipleTbCallback multipleTbCallback = new MultipleTbCallback(states.size(), msg.getCallback());
- states.forEach((cfId, state) -> cfStateService.removeState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), multipleTbCallback));
- ctx.stop(ctx.getSelf());
+ if (states.isEmpty()) {
+ msg.getCallback().onSuccess();
+ } else {
+ MultipleTbCallback multipleTbCallback = new MultipleTbCallback(states.size(), msg.getCallback());
+ states.forEach((cfId, state) -> cfStateService.removeState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), multipleTbCallback));
+ ctx.stop(ctx.getSelf());
+ }
} else {
var cfId = new CalculatedFieldId(msg.getEntityId().getId());
var state = states.remove(cfId);
if (state != null) {
cfStateService.removeState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), msg.getCallback());
+ } else {
+ msg.getCallback().onSuccess();
}
}
}
@@ -242,7 +248,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
justRestored = true;
}
if (state.isSizeOk()) {
- if (state.updateState(newArgValues) || justRestored) {
+ if (state.updateState(ctx, newArgValues) || justRestored) {
cfIdList = new ArrayList<>(cfIdList);
cfIdList.add(ctx.getCfId());
processStateIfReady(ctx, cfIdList, state, tbMsgId, tbMsgType, callback);
@@ -274,36 +280,34 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
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);
- boolean stateSizeOk;
- if (ctx.isInitialized() && state.isReady()) {
- try {
- CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(5, TimeUnit.SECONDS);
+ boolean stateSizeChecked = false;
+ try {
+ if (ctx.isInitialized() && state.isReady()) {
+ CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
state.checkStateSize(ctxId, ctx.getMaxStateSize());
- stateSizeOk = state.isSizeOk();
- if (stateSizeOk) {
+ stateSizeChecked = true;
+ if (state.isSizeOk()) {
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.getMaxStateSize());
- stateSizeOk = state.isSizeOk();
- if (stateSizeOk) {
- callback.onSuccess(); // State was updated but no calculation performed;
+ } catch (Exception e) {
+ throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).msgId(tbMsgId).msgType(tbMsgType).arguments(state.getArguments()).cause(e).build();
+ } finally {
+ if (!stateSizeChecked) {
+ state.checkStateSize(ctxId, ctx.getMaxStateSize());
+ }
+ if (state.isSizeOk()) {
+ cfStateService.persistState(ctxId, state, callback);
+ } else {
+ removeStateAndRaiseSizeException(ctxId, CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(), callback);
}
- }
- if (stateSizeOk) {
- cfStateService.persistState(ctxId, state, callback);
- } else {
- removeStateAndRaiseSizeException(ctxId, CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(), callback);
}
}
- private void removeStateAndRaiseSizeException(CalculatedFieldEntityCtxId ctxId, CalculatedFieldException ex, TbCallback callback) {
+ private void removeStateAndRaiseSizeException(CalculatedFieldEntityCtxId ctxId, CalculatedFieldException ex, TbCallback callback) throws CalculatedFieldException {
// We remove the state, but remember that it is over-sized in a local map.
cfStateService.removeState(ctxId, new TbCallback() {
@Override
@@ -316,6 +320,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
callback.onFailure(ex);
}
});
+ throw ex;
}
private Map mapToArguments(CalculatedFieldCtx ctx, List data) {
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 d26d89626e..1418821081 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
@@ -16,7 +16,6 @@
package org.thingsboard.server.actors.calculatedField;
import lombok.extern.slf4j.Slf4j;
-import org.thingsboard.common.util.DebugModeUtil;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorCtx;
import org.thingsboard.server.actors.TbActorRef;
@@ -51,8 +50,6 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CopyOnWriteArrayList;
import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto;
@@ -202,7 +199,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (fieldsCount > 0) {
MultipleTbCallback multiCallback = new MultipleTbCallback(fieldsCount, callback);
var entityId = msg.getEntityId();
- oldProfileCfs.forEach(ctx -> deleteCfForEntity(entityId, ctx.getCfId(), callback));
+ oldProfileCfs.forEach(ctx -> deleteCfForEntity(entityId, ctx.getCfId(), multiCallback));
newProfileCfs.forEach(ctx -> initCfForEntity(entityId, ctx, true, multiCallback));
} else {
callback.onSuccess();
@@ -309,7 +306,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (!entityIds.isEmpty()) {
//TODO: no need to do this if we cache all created actors and know which one belong to us;
var multiCallback = new MultipleTbCallback(entityIds.size(), callback);
- entityIds.forEach(id -> deleteCfForEntity(entityId, cfId, multiCallback));
+ entityIds.forEach(id -> deleteCfForEntity(id, cfId, multiCallback));
} else {
callback.onSuccess();
}
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
index d162b6a09a..ee6758df6a 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
@@ -28,8 +28,9 @@ import org.thingsboard.rule.engine.api.NotificationCenter;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.RuleEngineApiUsageStateService;
import org.thingsboard.rule.engine.api.RuleEngineAssetProfileCache;
+import org.thingsboard.rule.engine.api.RuleEngineCalculatedFieldQueueService;
import org.thingsboard.rule.engine.api.RuleEngineDeviceProfileCache;
-import org.thingsboard.rule.engine.api.RuleEngineDeviceStateManager;
+import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.RuleEngineRpcService;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.ScriptEngine;
@@ -725,7 +726,7 @@ public class DefaultTbContext implements TbContext {
}
@Override
- public RuleEngineDeviceStateManager getDeviceStateManager() {
+ public DeviceStateManager getDeviceStateManager() {
return mainCtx.getDeviceStateManager();
}
@@ -902,6 +903,11 @@ public class DefaultTbContext implements TbContext {
return mainCtx.getCalculatedFieldService();
}
+ @Override
+ public RuleEngineCalculatedFieldQueueService getCalculatedFieldQueueService() {
+ return mainCtx.getCalculatedFieldQueueService();
+ }
+
@Override
public boolean isExternalNodeForceAck() {
return mainCtx.isExternalNodeForceAck();
diff --git a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java
index fe85cf3f87..f899d0f480 100644
--- a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java
@@ -34,6 +34,8 @@ import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestController;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.script.api.tbel.TbelCfArg;
+import org.thingsboard.script.api.tbel.TbelCfCtx;
+import org.thingsboard.script.api.tbel.TbelCfSingleValueArg;
import org.thingsboard.script.api.tbel.TbelInvokeService;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EventInfo;
@@ -216,11 +218,14 @@ public class CalculatedFieldController extends BaseController {
@RequestBody JsonNode inputParams) {
String expression = inputParams.get("expression").asText();
Map arguments = Objects.requireNonNullElse(
- JacksonUtil.convertValue(inputParams.get("arguments"), new TypeReference