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 90709638e0..1a5e619d2c 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -41,6 +41,7 @@ import org.thingsboard.script.api.js.JsInvokeService; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.actors.service.ActorService; import org.thingsboard.server.actors.tenant.DebugTbRateLimits; +import org.thingsboard.server.cache.limits.RateLimitService; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.event.CalculatedFieldDebugEvent; import org.thingsboard.server.common.data.event.ErrorEvent; @@ -50,6 +51,7 @@ import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.limit.LimitedApi; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.TbActorMsg; @@ -107,8 +109,8 @@ 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.cache.CalculatedFieldEntityProfileCache; import org.thingsboard.server.service.cf.CalculatedFieldStateService; +import org.thingsboard.server.service.cf.cache.CalculatedFieldEntityProfileCache; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.component.ComponentDiscoveryService; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; @@ -130,6 +132,7 @@ import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.transport.TbCoreToTransportService; +import org.thingsboard.server.utils.DebugModeRateLimitsConfig; import java.io.PrintWriter; import java.io.StringWriter; @@ -182,8 +185,6 @@ public class ActorSystemContext { private final ConcurrentMap debugPerTenantLimits = new ConcurrentHashMap<>(); - private final ConcurrentMap cfDebugPerTenantLimits = new ConcurrentHashMap<>(); - public ConcurrentMap getDebugPerTenantLimits() { return debugPerTenantLimits; } @@ -452,6 +453,16 @@ public class ActorSystemContext { @Getter private ApiLimitService apiLimitService; + @Lazy + @Autowired(required = false) + @Getter + private RateLimitService rateLimitService; + + @Lazy + @Autowired(required = false) + @Getter + private DebugModeRateLimitsConfig debugModeRateLimitsConfig; + /** * The following Service will be null if we operate in tb-core mode */ @@ -609,14 +620,6 @@ public class ActorSystemContext { @Getter private long sessionReportTimeout; - @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.enabled:true}") - @Getter - private boolean debugPerTenantEnabled; - - @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.configuration:50000:3600}") - @Getter - private String debugPerTenantLimitsConfiguration; - @Value("${actors.rpc.submit_strategy:BURST}") @Getter private String rpcSubmitStrategy; @@ -641,14 +644,6 @@ public class ActorSystemContext { @Getter private String deviceStateNodeRateLimitConfig; - @Value("${actors.calculated_fields.debug_mode_rate_limits_per_tenant.enabled:true}") - @Getter - private boolean cfDebugPerTenantEnabled; - - @Value("${actors.calculated_fields.debug_mode_rate_limits_per_tenant.configuration:50000:3600}") - @Getter - private String cfDebugPerTenantLimitsConfiguration; - @Getter @Setter private TbActorSystem actorSystem; @@ -778,9 +773,9 @@ public class ActorSystemContext { } private boolean checkLimits(TenantId tenantId, TbMsg tbMsg, Throwable error) { - if (debugPerTenantEnabled) { + if (debugModeRateLimitsConfig.isRuleChainDebugPerTenantLimitsEnabled()) { DebugTbRateLimits debugTbRateLimits = debugPerTenantLimits.computeIfAbsent(tenantId, id -> - new DebugTbRateLimits(new TbRateLimits(debugPerTenantLimitsConfiguration), false)); + new DebugTbRateLimits(new TbRateLimits(debugModeRateLimitsConfig.getRuleChainDebugPerTenantLimitsConfiguration()), false)); if (!debugTbRateLimits.getTbRateLimits().tryConsume()) { if (!debugTbRateLimits.isRuleChainEventSaved()) { @@ -810,50 +805,51 @@ public class ActorSystemContext { Futures.addCallback(future, RULE_CHAIN_DEBUG_EVENT_ERROR_CALLBACK, MoreExecutors.directExecutor()); } - public void persistCalculatedFieldDebugEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, Map arguments, - UUID tbMsgId, TbMsgType tbMsgType, String result, Throwable error) { - if (cfDebugPerTenantEnabled) { - TbRateLimits rateLimits = cfDebugPerTenantLimits.computeIfAbsent(tenantId, id -> new TbRateLimits(cfDebugPerTenantLimitsConfiguration)); - - if (rateLimits.tryConsume()) { - try { - CalculatedFieldDebugEvent.CalculatedFieldDebugEventBuilder eventBuilder = CalculatedFieldDebugEvent.builder() - .tenantId(tenantId) - .entityId(calculatedFieldId.getId()) - .serviceId(getServiceId()) - .calculatedFieldId(calculatedFieldId) - .eventEntity(entityId); - if (tbMsgId != null) { - eventBuilder.msgId(tbMsgId); - } - if (tbMsgType != null) { - eventBuilder.msgType(tbMsgType.name()); - } - if (arguments != null) { - eventBuilder.arguments(JacksonUtil.toString( - arguments.entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toTbelCfArg())) - )); - } - if (result != null) { - eventBuilder.result(result); - } - if (error != null) { - eventBuilder.error(toString(error)); - } - - ListenableFuture future = eventService.saveAsync(eventBuilder.build()); - Futures.addCallback(future, CALCULATED_FIELD_DEBUG_EVENT_ERROR_CALLBACK, MoreExecutors.directExecutor()); - } catch (IllegalArgumentException ex) { - log.warn("Failed to persist calculated field debug message", ex); + public void persistCalculatedFieldDebugEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, Map arguments, UUID tbMsgId, TbMsgType tbMsgType, String result, Throwable error) { + if (checkLimits(tenantId)) { + try { + CalculatedFieldDebugEvent.CalculatedFieldDebugEventBuilder eventBuilder = CalculatedFieldDebugEvent.builder() + .tenantId(tenantId) + .entityId(calculatedFieldId.getId()) + .serviceId(getServiceId()) + .calculatedFieldId(calculatedFieldId) + .eventEntity(entityId); + if (tbMsgId != null) { + eventBuilder.msgId(tbMsgId); } - if (log.isTraceEnabled()) { - log.trace("[{}] Tenant level debug mode rate limit detected: {}", tenantId, calculatedFieldId); + if (tbMsgType != null) { + eventBuilder.msgType(tbMsgType.name()); + } + if (arguments != null) { + eventBuilder.arguments(JacksonUtil.toString( + arguments.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toTbelCfArg())) + )); + } + if (result != null) { + eventBuilder.result(result); } + if (error != null) { + eventBuilder.error(error.getMessage()); + } + + ListenableFuture future = eventService.saveAsync(eventBuilder.build()); + Futures.addCallback(future, CALCULATED_FIELD_DEBUG_EVENT_ERROR_CALLBACK, MoreExecutors.directExecutor()); + } catch (IllegalArgumentException ex) { + log.warn("Failed to persist calculated field debug message", ex); } } } + private boolean checkLimits(TenantId tenantId) { + if (debugModeRateLimitsConfig.isCalculatedFieldDebugPerTenantLimitsEnabled() && + !rateLimitService.checkRateLimit(LimitedApi.CALCULATED_FIELD_DEBUG_EVENTS, (Object) tenantId, debugModeRateLimitsConfig.getCalculatedFieldDebugPerTenantLimitsConfiguration())) { + log.trace("[{}] Calculated field debug event limits exceeded!", tenantId); + return false; + } + return true; + } + public static Exception toException(Throwable error) { return Exception.class.isInstance(error) ? (Exception) error : new Exception(error); } 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 c45c4c2ef0..34cf5695c1 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 @@ -38,8 +38,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelem import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldResult; -import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.CalculatedFieldStateService; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @@ -112,8 +112,16 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM log.info("Force reinitialization of CF: [{}].", cfCtx.getCfId()); states.remove(cfCtx.getCfId()); } - var cfState = getOrInitState(cfCtx); - processStateIfReady(cfCtx, Collections.singletonList(cfCtx.getCfId()), cfState, null, null, msg.getCallback()); + try { + var cfState = getOrInitState(cfCtx); + processStateIfReady(cfCtx, Collections.singletonList(cfCtx.getCfId()), cfState, null, null, msg.getCallback()); + } catch (Exception e) { + + if (DebugModeUtil.isDebugFailuresAvailable(cfCtx.getCalculatedField())) { + systemContext.persistCalculatedFieldDebugEvent(tenantId, cfCtx.getCfId(), entityId, null, null, null, null, e); + } + msg.getCallback().onFailure(e); + } } public void process(CalculatedFieldEntityDeleteMsg msg) { @@ -151,31 +159,45 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM var proto = msg.getProto(); var ctx = msg.getCtx(); var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback()); - List cfIds = getCalculatedFieldIds(proto); - if (cfIds.contains(ctx.getCfId())) { - callback.onSuccess(CALLBACKS_PER_CF); - } else { - if (proto.getTsDataCount() > 0) { - processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); - } else if (proto.getAttrDataCount() > 0) { - processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto)); - } else { + try { + List cfIds = getCalculatedFieldIds(proto); + if (cfIds.contains(ctx.getCfId())) { callback.onSuccess(CALLBACKS_PER_CF); + } else { + if (proto.getTsDataCount() > 0) { + processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); + } else if (proto.getAttrDataCount() > 0) { + processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto)); + } else { + callback.onSuccess(CALLBACKS_PER_CF); + } } + } catch (Exception e) { + if (DebugModeUtil.isDebugFailuresAvailable(ctx.getCalculatedField())) { + systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, null, null, null, null, e); + } + callback.onFailure(e); } } private void process(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, Collection cfIds, List cfIdList, MultipleTbCallback callback) { - if (cfIds.contains(ctx.getCfId())) { - callback.onSuccess(CALLBACKS_PER_CF); - } else { - if (proto.getTsDataCount() > 0) { - processTelemetry(ctx, proto, cfIdList, callback); - } else if (proto.getAttrDataCount() > 0) { - processAttributes(ctx, proto, cfIdList, callback); - } else { + try { + if (cfIds.contains(ctx.getCfId())) { callback.onSuccess(CALLBACKS_PER_CF); + } else { + if (proto.getTsDataCount() > 0) { + processTelemetry(ctx, proto, cfIdList, callback); + } else if (proto.getAttrDataCount() > 0) { + processAttributes(ctx, proto, cfIdList, callback); + } else { + callback.onSuccess(CALLBACKS_PER_CF); + } + } + } catch (Exception e) { + if (DebugModeUtil.isDebugFailuresAvailable(ctx.getCalculatedField())) { + systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, null, null, null, null, e); } + callback.onFailure(e); } } @@ -230,22 +252,18 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM @SneakyThrows private void processStateIfReady(CalculatedFieldCtx ctx, List cfIdList, CalculatedFieldState state, UUID tbMsgId, TbMsgType tbMsgType, TbCallback callback) { + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId); if (state.isReady() && ctx.isInitialized()) { - try { - CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(5, TimeUnit.SECONDS); - 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.getResultMap()), null); - } - } catch (Exception e) { - if (DebugModeUtil.isDebugFailuresAvailable(ctx.getCalculatedField())) { - systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, null, e); - } + 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); } } else { callback.onSuccess(); // State was updated but no calculation performed; } - cfStateService.persistState(ctx, new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), state, callback); + cfStateService.persistState(ctxId, state, callback); } 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 22492f3c4f..c3aa378524 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 @@ -41,9 +41,9 @@ 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; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; -import org.thingsboard.server.service.cf.CalculatedFieldStateService; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; @@ -53,7 +53,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -263,6 +262,13 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware callback.onSuccess(); } else { var newCfCtx = new CalculatedFieldCtx(newCf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService()); + try { + newCfCtx.init(); + } catch (Exception e) { + if (DebugModeUtil.isDebugAllAvailable(newCf)) { + systemContext.persistCalculatedFieldDebugEvent(newCf.getTenantId(), newCf.getId(), newCf.getEntityId(), null, null, null, null, e); + } + } calculatedFields.put(newCf.getId(), newCfCtx); List oldCfList = entityIdCalculatedFields.get(newCf.getEntityId()); List newCfList = new ArrayList<>(oldCfList.size()); @@ -287,13 +293,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware // Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead) var stateChanges = newCfCtx.hasStateChanges(oldCfCtx); if (stateChanges || newCfCtx.hasOtherSignificantChanges(oldCfCtx)) { - try { - newCfCtx.init(); - } catch (Exception e) { - if (DebugModeUtil.isDebugAllAvailable(newCf)) { - systemContext.persistCalculatedFieldDebugEvent(newCf.getTenantId(), newCf.getId(), newCf.getEntityId(), null, null, null, null, e); - } - } initCf(newCfCtx, callback, stateChanges); } else { callback.onSuccess(); 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 5c8b40dcad..39f7fc91a8 100644 --- a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java +++ b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java @@ -33,6 +33,7 @@ import org.springframework.web.bind.annotation.ResponseBody; 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.TbelInvokeService; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -84,14 +85,24 @@ public class CalculatedFieldController extends BaseController { private static final String TEST_SCRIPT_EXPRESSION = "Execute the Script expression and return the result. The format of request: \n\n" + MARKDOWN_CODE_BLOCK_START + "{\n" + - " \"expression\": \"var temp = 0; foreach(element: temperature.entrySet()) { temp += element.getValue(); } var avgTemperature = temp / temperature.size(); var adjustedTemperature = avgTemperature + 0.1 * humidity; return { \\\"adjustedTemperature\\\": adjustedTemperature };\",\n" + + " \"expression\": \"var temp = 0; foreach(element: temperature.values) {temp += element.value;} var avgTemperature = temp / temperature.values.size(); var adjustedTemperature = avgTemperature + 0.1 * humidity.value; return {\\\"adjustedTemperature\\\": adjustedTemperature};\",\n" + " \"arguments\": {\n" + " \"temperature\": {\n" + - " \"14327856345\": 22.4,\n" + - " \"14327857298\": 21.9,\n" + - " \"14327857510\": 22.0\n" + + " \"type\": \"TS_ROLLING\",\n" + + " \"timeWindow\": {\n" + + " \"startTs\": 1739775630002,\n" + + " \"endTs\": 65432211,\n" + + " \"limit\": 5\n" + + " },\n" + + " \"values\": [\n" + + " { \"ts\": 1739775639851, \"value\": 23 },\n" + + " { \"ts\": 1739775664561, \"value\": 43 },\n" + + " { \"ts\": 1739775713079, \"value\": 15 },\n" + + " { \"ts\": 1739775999522, \"value\": 34 },\n" + + " { \"ts\": 1739776228452, \"value\": 22 }\n" + + " ]\n" + " },\n" + - " \"humidity\": 42\n" + + " \"humidity\": { \"type\": \"SINGLE_VALUE\", \"ts\": 1739776478057, \"value\": 23 }\n" + " }\n" + "}" + MARKDOWN_CODE_BLOCK_END @@ -172,8 +183,8 @@ public class CalculatedFieldController extends BaseController { @io.swagger.v3.oas.annotations.parameters.RequestBody(description = "Test calculated field TBEL expression.") @RequestBody JsonNode inputParams) { String expression = inputParams.get("expression").asText(); - Map arguments = Objects.requireNonNullElse( - JacksonUtil.convertValue(inputParams.get("arguments"), new TypeReference>() { + Map arguments = Objects.requireNonNullElse( + JacksonUtil.convertValue(inputParams.get("arguments"), new TypeReference>() { }), Collections.emptyMap() ); diff --git a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java index 8b958df68b..3fbedcebaf 100644 --- a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java +++ b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java @@ -35,8 +35,8 @@ import org.thingsboard.server.common.data.SystemParams; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.mobile.qrCodeSettings.QrCodeSettings; import org.thingsboard.server.common.data.mobile.qrCodeSettings.QRCodeConfig; +import org.thingsboard.server.common.data.mobile.qrCodeSettings.QrCodeSettings; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.settings.UserSettings; import org.thingsboard.server.common.data.settings.UserSettingsType; @@ -46,6 +46,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.UserPrincipal; import org.thingsboard.server.service.sync.vc.EntitiesVersionControlService; +import org.thingsboard.server.utils.DebugModeRateLimitsConfig; import java.util.Collections; import java.util.List; @@ -74,12 +75,6 @@ public class SystemInfoController extends BaseController { @Value("${debug.settings.default_duration:15}") private int defaultDebugDurationMinutes; - @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.enabled:true}") - private boolean ruleChainDebugPerTenantLimitsEnabled; - - @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.configuration:50000:3600}") - private String ruleChainDebugPerTenantLimitsConfiguration; - @Autowired(required = false) private BuildProperties buildProperties; @@ -89,6 +84,9 @@ public class SystemInfoController extends BaseController { @Autowired private QrCodeSettingService qrCodeSettingService; + @Autowired + private DebugModeRateLimitsConfig debugModeRateLimitsConfig; + @PostConstruct public void init() { JsonNode info = buildInfoObject(); @@ -152,8 +150,11 @@ public class SystemInfoController extends BaseController { DefaultTenantProfileConfiguration tenantProfileConfiguration = tenantProfileCache.get(tenantId).getDefaultProfileConfiguration(); systemParams.setMaxResourceSize(tenantProfileConfiguration.getMaxResourceSize()); systemParams.setMaxDebugModeDurationMinutes(DebugModeUtil.getMaxDebugAllDuration(tenantProfileConfiguration.getMaxDebugModeDurationMinutes(), defaultDebugDurationMinutes)); - if (ruleChainDebugPerTenantLimitsEnabled) { - systemParams.setRuleChainDebugPerTenantLimitsConfiguration(ruleChainDebugPerTenantLimitsConfiguration); + if (debugModeRateLimitsConfig.isRuleChainDebugPerTenantLimitsEnabled()) { + systemParams.setRuleChainDebugPerTenantLimitsConfiguration(debugModeRateLimitsConfig.getRuleChainDebugPerTenantLimitsConfiguration()); + } + if (debugModeRateLimitsConfig.isCalculatedFieldDebugPerTenantLimitsEnabled()) { + systemParams.setCalculatedFieldDebugPerTenantLimitsConfiguration(debugModeRateLimitsConfig.getCalculatedFieldDebugPerTenantLimitsConfiguration()); } } systemParams.setMobileQrEnabled(Optional.ofNullable(qrCodeSettingService.findQrCodeSettings(TenantId.SYS_TENANT_ID)) diff --git a/application/src/main/java/org/thingsboard/server/exception/CalculatedFieldStateException.java b/application/src/main/java/org/thingsboard/server/exception/CalculatedFieldStateException.java new file mode 100644 index 0000000000..50ab512cb9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/exception/CalculatedFieldStateException.java @@ -0,0 +1,24 @@ +/** + * 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.server.exception; + +public class CalculatedFieldStateException extends RuntimeException { + + public CalculatedFieldStateException(String message) { + super(message); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java index 0e5fc93c6e..e390d9e55c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java @@ -18,32 +18,14 @@ package org.thingsboard.server.service.cf; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.calculatedField.CalculatedFieldStateRestoreMsg; -import org.thingsboard.server.common.data.StringUtils; -import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.id.CalculatedFieldId; -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.data.kv.BasicKvEntry; import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.common.util.KvProtoUtil; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityCtxIdProto; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.exception.CalculatedFieldStateException; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; -import org.thingsboard.server.gen.transport.TransportProtos.TsDoubleValProto; -import org.thingsboard.server.gen.transport.TransportProtos.TsRollingArgumentProto; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; -import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; -import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; -import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; -import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; -import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry; -import java.util.Optional; -import java.util.TreeMap; -import java.util.UUID; +import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; +import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; public abstract class AbstractCalculatedFieldStateService implements CalculatedFieldStateService { @@ -51,15 +33,14 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF private ActorSystemContext actorSystemContext; @Override - public final void persistState(CalculatedFieldCtx ctx, CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) { - CalculatedFieldStateMsgProto stateMsg = toProto(stateId, state); - long maxStateSizeInKBytes = ctx.getMaxStateSizeInKBytes(); - if (maxStateSizeInKBytes <= 0 || stateMsg.getSerializedSize() <= maxStateSizeInKBytes) { - doPersist(stateId, stateMsg, callback); + public final void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) { + if (state.isStateTooLarge()) { + throw new CalculatedFieldStateException("State size exceeds the maximum allowed limit. The state will not be persisted to RocksDB."); } + doPersist(stateId, toProto(stateId, state), callback); } - protected abstract void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateMsgProto stateMsgProto, TbCallback callback); + protected abstract void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback); @Override public final void removeState(CalculatedFieldEntityCtxId stateId, TbCallback callback) { @@ -68,110 +49,12 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF protected abstract void doRemove(CalculatedFieldEntityCtxId stateId, TbCallback callback); - protected void processRestoredState(CalculatedFieldStateMsgProto stateMsg) { - CalculatedFieldEntityCtxId stateId = fromProto(stateMsg.getId()); - CalculatedFieldState state = stateMsg.hasState() ? fromProto(stateMsg.getState()) : null; - actorSystemContext.tell(new CalculatedFieldStateRestoreMsg(stateId, state)); + protected void processRestoredState(CalculatedFieldStateProto stateMsg) { + var id = fromProto(stateMsg.getId()); + var state = fromProto(stateMsg); + actorSystemContext.tell(new CalculatedFieldStateRestoreMsg(id, state)); } - protected CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) { - return CalculatedFieldEntityCtxIdProto.newBuilder() - .setTenantIdMSB(ctxId.tenantId().getId().getMostSignificantBits()) - .setTenantIdLSB(ctxId.tenantId().getId().getLeastSignificantBits()) - .setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits()) - .setEntityType(ctxId.entityId().getEntityType().name()) - .setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits()) - .setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits()) - .build(); - } - - protected CalculatedFieldEntityCtxId fromProto(CalculatedFieldEntityCtxIdProto ctxIdProto) { - TenantId tenantId = TenantId.fromUUID(new UUID(ctxIdProto.getTenantIdMSB(), ctxIdProto.getTenantIdLSB())); - 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); - } - - protected CalculatedFieldStateMsgProto toProto(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state) { - var stateProto = CalculatedFieldStateProto.newBuilder() - .setType(state.getType().name()); - state.getArguments().forEach((argName, argEntry) -> { - if (argEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - stateProto.addSingleValueArguments(toSingleValueArgumentProto(argName, singleValueArgumentEntry)); - } else if (argEntry instanceof TsRollingArgumentEntry rollingArgumentEntry) { - stateProto.addRollingValueArguments(toRollingArgumentProto(argName, rollingArgumentEntry)); - } - }); - return CalculatedFieldStateMsgProto.newBuilder() - .setId(toProto(stateId)) - .setState(stateProto) - .build(); - } - - private TransportProtos.SingleValueArgumentProto toSingleValueArgumentProto(String argName, SingleValueArgumentEntry entry) { - TransportProtos.SingleValueArgumentProto.Builder builder = TransportProtos.SingleValueArgumentProto.newBuilder() - .setArgName(argName); - - if (entry != SingleValueArgumentEntry.EMPTY) { - builder.setValue(KvProtoUtil.toTsValueProto(entry.getTs(), entry.getKvEntryValue())); - } - - Optional.ofNullable(entry.getVersion()).ifPresent(builder::setVersion); - - return builder.build(); - } - - private TsRollingArgumentProto toRollingArgumentProto(String argName, TsRollingArgumentEntry entry) { - TsRollingArgumentProto.Builder builder = TsRollingArgumentProto.newBuilder().setKey(argName); - - if (entry != TsRollingArgumentEntry.EMPTY) { - entry.getTsRecords().forEach((ts, value) -> builder.addTsValue(TsDoubleValProto.newBuilder().setTs(ts).setValue(value).build())); - } - - return builder.build(); - } - - private CalculatedFieldState fromProto(CalculatedFieldStateProto proto) { - if (StringUtils.isEmpty(proto.getType())) { - return null; - } - - CalculatedFieldType type = CalculatedFieldType.valueOf(proto.getType()); - - CalculatedFieldState state = switch (type) { - case SIMPLE -> new SimpleCalculatedFieldState(); - case SCRIPT -> new ScriptCalculatedFieldState(); - }; - - proto.getSingleValueArgumentsList().forEach(argProto -> - state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto))); - - if (CalculatedFieldType.SCRIPT.equals(type)) { - proto.getRollingValueArgumentsList().forEach(argProto -> - state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto))); - } - return state; - } - - private SingleValueArgumentEntry fromSingleValueArgumentProto(TransportProtos.SingleValueArgumentProto proto) { - if (!proto.hasValue()) { - return (SingleValueArgumentEntry) SingleValueArgumentEntry.EMPTY; - } - TransportProtos.TsValueProto tsValueProto = proto.getValue(); - long ts = tsValueProto.getTs(); - BasicKvEntry kvEntry = (BasicKvEntry) KvProtoUtil.fromTsValueProto(proto.getArgName(), tsValueProto); - return new SingleValueArgumentEntry(ts, kvEntry, proto.getVersion()); - } - - private TsRollingArgumentEntry fromRollingArgumentProto(TsRollingArgumentProto proto) { - if (proto.getTsValueCount() <= 0) { - return (TsRollingArgumentEntry) TsRollingArgumentEntry.EMPTY; - } - TreeMap tsRecords = new TreeMap<>(); - proto.getTsValueList().forEach(tsValueProto -> tsRecords.put(tsValueProto.getTs(), tsValueProto.getValue())); - return new TsRollingArgumentEntry(tsRecords); - } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index 52e0a0151b..13ac4972e2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -15,22 +15,16 @@ */ package org.thingsboard.server.service.cf; +import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.configuration.OutputType; -import java.util.Map; - @Data public final class CalculatedFieldResult { private final OutputType type; private final AttributeScope scope; - private final Map resultMap; + private final JsonNode result; - public CalculatedFieldResult(OutputType type, AttributeScope scope, Map resultMap) { - this.type = type; - this.scope = scope; - this.resultMap = resultMap == null ? Map.of() : Map.copyOf(resultMap); - } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java index 05f9dab36b..df2b04c20a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java @@ -17,15 +17,15 @@ package org.thingsboard.server.service.cf; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.exception.CalculatedFieldStateException; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; -import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; import java.util.Set; public interface CalculatedFieldStateService { - void persistState(CalculatedFieldCtx ctx, CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback); + void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException; void removeState(CalculatedFieldEntityCtxId stateId, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index f2c6976cbf..3abaea75b4 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -15,8 +15,6 @@ */ package org.thingsboard.server.service.cf; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; @@ -60,8 +58,6 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityCtxIdProto; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldIdProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto.Builder; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; @@ -92,6 +88,7 @@ import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.SCOPE; +import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; @TbRuleEngineComponent @Service @@ -156,8 +153,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP OutputType type = calculatedFieldResult.getType(); TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; - ObjectNode payload = createJsonPayload(calculatedFieldResult); - TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(JacksonUtil.writeValueAsString(payload)).build(); + TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(JacksonUtil.writeValueAsString(calculatedFieldResult.getResult())).build(); clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { @Override public void onSuccess(TbQueueMsgMetadata metadata) { @@ -229,19 +225,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP return builder.build(); } - //TODO: IM: move to utils; - private CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) { - return CalculatedFieldEntityCtxIdProto.newBuilder() - .setTenantIdMSB(ctxId.tenantId().getId().getMostSignificantBits()) - .setTenantIdLSB(ctxId.tenantId().getId().getLeastSignificantBits()) - .setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits()) - .setEntityType(ctxId.entityId().getEntityType().name()) - .setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits()) - .setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits()) - .build(); - } - private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { return switch (argument.getRefEntityKey().getType()) { case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); @@ -264,7 +247,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP if (kvEntry.isPresent() && kvEntry.get().getValue() != null) { return ArgumentEntry.createSingleValueArgument(kvEntry.get()); } else { - return SingleValueArgumentEntry.EMPTY; + return new SingleValueArgumentEntry(); } }, calculatedFieldCallbackExecutor); } @@ -279,7 +262,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); - return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); + return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? new TsRollingArgumentEntry(limit, timeWindow) : ArgumentEntry.createTsRollingArgument(tsRolling, limit, timeWindow), calculatedFieldCallbackExecutor); } private KvEntry createDefaultKvEntry(Argument argument) { @@ -297,13 +280,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP return new StringDataEntry(key, defaultValue); } - private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { - ObjectNode payload = JacksonUtil.newObjectNode(); - Map resultMap = calculatedFieldResult.getResultMap(); - resultMap.forEach((k, v) -> payload.set(k, JacksonUtil.convertValue(v, JsonNode.class))); - return payload; - } - private CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) { return switch (ctx.getCfType()) { case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames()); @@ -311,13 +287,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP }; } - private CalculatedFieldIdProto toProto(CalculatedFieldId cfId) { - return CalculatedFieldIdProto.newBuilder() - .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) - .build(); - } - private static class TbCallbackWrapper implements TbQueueCallback { private final TbCallback callback; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java index 74e2d73601..13656d5b0a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java @@ -42,17 +42,16 @@ public interface ArgumentEntry { boolean updateEntry(ArgumentEntry entry); + boolean isEmpty(); + static ArgumentEntry createSingleValueArgument(KvEntry kvEntry) { return new SingleValueArgumentEntry(kvEntry); } - static ArgumentEntry createTsRollingArgument(List kvEntries) { - return new TsRollingArgumentEntry(kvEntries); + static ArgumentEntry createTsRollingArgument(List kvEntries, int limit, long timeWindow) { + return new TsRollingArgumentEntry(kvEntries, limit, timeWindow); } - @JsonIgnore - ArgumentEntry copy(); - TbelCfArg toTbelCfArg(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index f1ce8038c7..5ea4707ff3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -17,6 +17,8 @@ package org.thingsboard.server.service.cf.ctx.state; import lombok.AllArgsConstructor; import lombok.Data; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.utils.CalculatedFieldUtils; import java.util.ArrayList; import java.util.HashMap; @@ -29,6 +31,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected List requiredArguments; protected Map arguments; + protected boolean stateTooLarge; public BaseCalculatedFieldState(List requiredArguments) { this.requiredArguments = requiredArguments; @@ -36,7 +39,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { } public BaseCalculatedFieldState() { - this(new ArrayList<>(), new HashMap<>()); + this(new ArrayList<>(), new HashMap<>(), false); } @Override @@ -52,7 +55,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { ArgumentEntry newEntry = entry.getValue(); ArgumentEntry existingEntry = arguments.get(key); - if (existingEntry == null || existingEntry == SingleValueArgumentEntry.EMPTY || existingEntry == TsRollingArgumentEntry.EMPTY) { + if (existingEntry == null) { validateNewEntry(newEntry); arguments.put(key, newEntry); stateUpdated = true; @@ -67,8 +70,14 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { @Override public boolean isReady() { return arguments.keySet().containsAll(requiredArguments) && - !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && - !arguments.containsValue(TsRollingArgumentEntry.EMPTY); + arguments.values().stream().noneMatch(ArgumentEntry::isEmpty); + } + + @Override + public void checkStateSize(CalculatedFieldEntityCtxId ctxId, long maxStateSize) { + if (maxStateSize > 0 && CalculatedFieldUtils.toProto(ctxId, this).getSerializedSize() > maxStateSize) { + setStateTooLarge(true); + } } protected abstract void validateNewEntry(ArgumentEntry newEntry); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java index 6bea6ce705..b54ef744a0 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java @@ -18,18 +18,12 @@ package org.thingsboard.server.service.cf.ctx.state; import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.ListenableFuture; -import java.util.Map; - public interface CalculatedFieldScriptEngine { ListenableFuture executeScriptAsync(Object[] args); - ListenableFuture> executeToMapAsync(Object[] args); - ListenableFuture executeJsonAsync(Object[] args); - ListenableFuture> executeToMapTransform(Object result); - void destroy(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index 0792b8240d..0219dd2593 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.service.cf.CalculatedFieldResult; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import java.util.List; import java.util.Map; @@ -41,8 +42,6 @@ public interface CalculatedFieldState { Map getArguments(); - List getRequiredArguments(); - void setRequiredArguments(List requiredArguments); boolean updateState(Map argumentValues); @@ -51,4 +50,9 @@ public interface CalculatedFieldState { @JsonIgnore boolean isReady(); + + boolean isStateTooLarge(); + + void checkStateSize(CalculatedFieldEntityCtxId ctxId, long maxStateSize); + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java index 6dcbf33067..c298f9cbe3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java @@ -19,7 +19,6 @@ import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import io.jsonwebtoken.lang.Collections; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.ScriptType; @@ -27,7 +26,6 @@ import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.common.data.id.TenantId; import javax.script.ScriptException; -import java.util.Map; import java.util.UUID; import java.util.concurrent.ExecutionException; @@ -72,24 +70,11 @@ public class CalculatedFieldTbelScriptEngine implements CalculatedFieldScriptEng }, MoreExecutors.directExecutor()); } - @Override - public ListenableFuture> executeToMapAsync(Object[] args) { - return Futures.transformAsync(executeScriptAsync(args), this::executeToMapTransform, MoreExecutors.directExecutor()); - } - @Override public ListenableFuture executeJsonAsync(Object[] args) { return Futures.transform(executeScriptAsync(args), JacksonUtil::valueToTree, MoreExecutors.directExecutor()); } - @Override - public ListenableFuture> executeToMapTransform(Object result) { - if (result instanceof Map) { - return Futures.immediateFuture((Map) result); - } - throw new IllegalArgumentException("Wrong result type: [" + result.getClass().getName() + "]"); - } - @Override public void destroy() { tbelInvokeService.release(this.scriptId); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java index 152b1affd9..cce7e5ef22 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java @@ -26,7 +26,7 @@ import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueMsgMetadata; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -60,8 +60,8 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta @Value("${queue.calculated_fields.consumer_per_partition:true}") private boolean consumerPerPartition; - private MainQueueConsumerManager, CalculatedFieldQueueConfig> stateConsumer; - private TbKafkaProducerTemplate> stateProducer; + private MainQueueConsumerManager, CalculatedFieldQueueConfig> stateConsumer; + private TbKafkaProducerTemplate> stateProducer; protected ExecutorService consumersExecutor; protected ExecutorService mgmtExecutor; @@ -75,11 +75,11 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta this.mgmtExecutor = ThingsBoardExecutors.newWorkStealingPool(Math.max(Runtime.getRuntime().availableProcessors(), 4), "cf-state-mgmt"); this.scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("cf-state-consumer-scheduler"); - this.stateConsumer = MainQueueConsumerManager., CalculatedFieldQueueConfig>builder() + this.stateConsumer = MainQueueConsumerManager., CalculatedFieldQueueConfig>builder() .queueKey(QueueKey.CF_STATES) .config(CalculatedFieldQueueConfig.of(consumerPerPartition, (int) pollInterval)) .msgPackProcessor((msgs, consumer, config) -> { - for (TbProtoQueueMsg msg : msgs) { + for (TbProtoQueueMsg msg : msgs) { try { processRestoredState(msg.getValue()); } catch (Throwable t) { @@ -97,11 +97,11 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta .scheduler(scheduler) .taskExecutor(mgmtExecutor) .build(); - this.stateProducer = (TbKafkaProducerTemplate>) queueFactory.createCalculatedFieldStateProducer(); + this.stateProducer = (TbKafkaProducerTemplate>) queueFactory.createCalculatedFieldStateProducer(); } @Override - protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateMsgProto stateMsgProto, TbCallback callback) { + protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback) { TopicPartitionInfo tpi = partitionService.resolve(QueueKey.CF_STATES, stateId.entityId()); stateProducer.send(tpi, stateId.toKey(), new TbProtoQueueMsg<>(stateId.entityId().getId(), stateMsgProto), new TbQueueCallback() { @Override @@ -122,9 +122,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta @Override protected void doRemove(CalculatedFieldEntityCtxId stateId, TbCallback callback) { - doPersist(stateId, CalculatedFieldStateMsgProto.newBuilder() - .setId(toProto(stateId)) - .build(), callback); + //TODO: vklimov } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java index eb7f4818d1..61e752b9d1 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java @@ -22,7 +22,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; import org.thingsboard.server.service.cf.CfRocksDb; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; @@ -40,7 +40,7 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS private Set partitions; @Override - protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateMsgProto stateMsgProto, TbCallback callback) { + protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback) { cfRocksDb.put(stateId.toKey(), stateMsgProto.toByteArray()); callback.onSuccess(); } @@ -61,7 +61,7 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS cfRocksDb.forEach((key, value) -> { try { - processRestoredState(CalculatedFieldStateMsgProto.parseFrom(value)); + processRestoredState(CalculatedFieldStateProto.parseFrom(value)); } catch (InvalidProtocolBufferException e) { log.error("[{}] Failed to process restored state", key, e); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index d45b5eaa27..94f0894ff3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cf.ctx.state; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -22,18 +23,11 @@ import lombok.Data; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.thingsboard.script.api.tbel.TbelCfArg; -import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; -import org.thingsboard.script.api.tbel.TbelCfTsDoubleVal; -import org.thingsboard.script.api.tbel.TbelCfTsRollingArg; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.service.cf.CalculatedFieldResult; -import java.util.ArrayList; import java.util.List; -import java.util.Map; -import java.util.TreeMap; @Data @Slf4j @@ -55,21 +49,11 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { - arguments.forEach((key, argumentEntry) -> { - if (argumentEntry instanceof TsRollingArgumentEntry tsRollingEntry) { - Argument argument = ctx.getArguments().get(key); - TreeMap tsRecords = tsRollingEntry.getTsRecords(); - if (tsRecords.size() > argument.getLimit()) { - tsRecords.pollFirstEntry(); - } - tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); - } - }); Object[] args = ctx.getArgNames().stream() .map(this::toTbelArgument) .toArray(); - ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); + ListenableFuture resultFuture = ctx.getCalculatedFieldScriptEngine().executeJsonAsync(args); Output output = ctx.getOutput(); return Futures.transform(resultFuture, result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 01cf9cff70..fa76a63e25 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.Data; import lombok.NoArgsConstructor; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.kv.BasicKvEntry; @@ -63,7 +64,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { double expressionResult = expr.evaluate(); Output output = ctx.getOutput(); - return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); + return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), JacksonUtil.valueToTree(Map.of(output.getName(), expressionResult)))); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java index 5614e592d7..c096233224 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java @@ -34,8 +34,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; @AllArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { - public static final ArgumentEntry EMPTY = new SingleValueArgumentEntry(0); - private long ts; private BasicKvEntry kvEntryValue; private Long version; @@ -63,27 +61,19 @@ public class SingleValueArgumentEntry implements ArgumentEntry { this.kvEntryValue = ProtoUtils.basicKvEntryFromKvEntry(entry); } - /** - * Internal constructor to create immutable SingleValueArgumentEntry.EMPTY - * */ - private SingleValueArgumentEntry(int ignored) { - this.ts = System.currentTimeMillis(); - this.kvEntryValue = null; - } - @Override public ArgumentEntryType getType() { return ArgumentEntryType.SINGLE_VALUE; } - @JsonIgnore - public Object getValue() { - return kvEntryValue.getValue(); + @Override + public boolean isEmpty() { + return kvEntryValue == null; } - @Override - public ArgumentEntry copy() { - return new SingleValueArgumentEntry(this.ts, this.kvEntryValue, this.version); + @JsonIgnore + public Object getValue() { + return isEmpty() ? null : kvEntryValue.getValue(); } @Override @@ -105,7 +95,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { // TODO: should we persist updated ts and version values? BasicKvEntry newValue = singleValueEntry.getKvEntryValue(); - if (this.kvEntryValue.getValue().equals(newValue.getValue())) { + if (this.kvEntryValue != null && this.kvEntryValue.getValue().equals(newValue.getValue())) { return false; } this.kvEntryValue = singleValueEntry.getKvEntryValue(); 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 f7bc891c41..866e2f5e09 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 @@ -20,21 +20,17 @@ import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.math.NumberUtils; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfTsDoubleVal; import org.thingsboard.script.api.tbel.TbelCfTsRollingArg; -import org.thingsboard.server.common.data.kv.BasicKvEntry; -import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; -import org.thingsboard.server.common.util.ProtoUtils; +import org.thingsboard.server.exception.CalculatedFieldStateException; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.TreeMap; -import java.util.stream.Collectors; @Data @NoArgsConstructor @@ -42,21 +38,26 @@ import java.util.stream.Collectors; @Slf4j public class TsRollingArgumentEntry implements ArgumentEntry { - public static final ArgumentEntry EMPTY = new TsRollingArgumentEntry(0); - - private static final int MAX_ROLLING_ARGUMENT_ENTRY_SIZE = 1000; - + private Integer limit; + private Long timeWindow; private TreeMap tsRecords = new TreeMap<>(); - public TsRollingArgumentEntry(List kvEntries) { + public TsRollingArgumentEntry(List kvEntries, int limit, long timeWindow) { + this.limit = limit; + this.timeWindow = timeWindow; kvEntries.forEach(tsKvEntry -> addTsRecord(tsKvEntry.getTs(), tsKvEntry)); } - /** - * Internal constructor to create immutable TsRollingArgumentEntry.EMPTY - */ - private TsRollingArgumentEntry(int ignored) { + public TsRollingArgumentEntry(TreeMap tsRecords, int limit, long timeWindow) { + this.tsRecords = tsRecords; + this.limit = limit; + this.timeWindow = timeWindow; + } + + public TsRollingArgumentEntry(int limit, long timeWindow) { this.tsRecords = new TreeMap<>(); + this.limit = limit; + this.timeWindow = timeWindow; } @Override @@ -64,15 +65,15 @@ public class TsRollingArgumentEntry implements ArgumentEntry { return ArgumentEntryType.TS_ROLLING; } - @JsonIgnore @Override - public Object getValue() { - return tsRecords; + public boolean isEmpty() { + return tsRecords.isEmpty(); } + @JsonIgnore @Override - public ArgumentEntry copy() { - return new TsRollingArgumentEntry(new TreeMap<>(tsRecords)); + public Object getValue() { + return tsRecords; } @Override @@ -81,11 +82,11 @@ public class TsRollingArgumentEntry implements ArgumentEntry { for (var e : tsRecords.entrySet()) { values.add(new TbelCfTsDoubleVal(e.getKey(), e.getValue())); } - return new TbelCfTsRollingArg(values); + return new TbelCfTsRollingArg(limit, timeWindow, values); } @Override - public boolean updateEntry(ArgumentEntry entry) { + public boolean updateEntry(ArgumentEntry entry) throws CalculatedFieldStateException { if (entry instanceof TsRollingArgumentEntry tsRollingEntry) { updateTsRollingEntry(tsRollingEntry); } else if (entry instanceof SingleValueArgumentEntry singleValueEntry) { @@ -107,26 +108,31 @@ public class TsRollingArgumentEntry implements ArgumentEntry { } private void addTsRecord(Long ts, KvEntry value) { - switch (value.getDataType()) { - case LONG -> value.getLongValue().ifPresent(aLong -> tsRecords.put(ts, aLong.doubleValue())); - case DOUBLE -> value.getDoubleValue().ifPresent(aDouble -> tsRecords.put(ts, aDouble)); - case BOOLEAN -> value.getBooleanValue().ifPresent(aBoolean -> tsRecords.put(ts, aBoolean ? 1.0 : 0.0)); - case STRING -> value.getStrValue().ifPresent(aString -> tsRecords.put(ts, Double.parseDouble(aString))); - case JSON -> value.getJsonValue().ifPresent(aString -> tsRecords.put(ts, Double.parseDouble(aString))); - //TODO: try catch + try { + switch (value.getDataType()) { + case LONG -> value.getLongValue().ifPresent(aLong -> tsRecords.put(ts, aLong.doubleValue())); + case DOUBLE -> value.getDoubleValue().ifPresent(aDouble -> tsRecords.put(ts, aDouble)); + case BOOLEAN -> value.getBooleanValue().ifPresent(aBoolean -> tsRecords.put(ts, aBoolean ? 1.0 : 0.0)); + 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) { + log.warn("Time series rolling arguments supports only numeric values."); +// throw new IllegalArgumentException("Time series rolling arguments supports only numeric values."); } - pollFirstEntryIfNeeded(); } private void addTsRecord(Long ts, double value) { tsRecords.put(ts, value); - pollFirstEntryIfNeeded(); + cleanupExpiredRecords(); } - private void pollFirstEntryIfNeeded() { - if (tsRecords.size() > MAX_ROLLING_ARGUMENT_ENTRY_SIZE) { + private void cleanupExpiredRecords() { + if (tsRecords.size() > limit) { tsRecords.pollFirstEntry(); } + tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - timeWindow); } } diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java new file mode 100644 index 0000000000..112db1a659 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -0,0 +1,144 @@ +/** + * 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.server.utils; + +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +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.data.kv.BasicKvEntry; +import org.thingsboard.server.common.util.KvProtoUtil; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityCtxIdProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; +import org.thingsboard.server.gen.transport.TransportProtos.SingleValueArgumentProto; +import org.thingsboard.server.gen.transport.TransportProtos.TsDoubleValProto; +import org.thingsboard.server.gen.transport.TransportProtos.TsRollingArgumentProto; +import org.thingsboard.server.gen.transport.TransportProtos.TsValueProto; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; +import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry; + +import java.util.Optional; +import java.util.TreeMap; +import java.util.UUID; + +public class CalculatedFieldUtils { + + public static CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) { + return CalculatedFieldEntityCtxIdProto.newBuilder() + .setTenantIdMSB(ctxId.tenantId().getId().getMostSignificantBits()) + .setTenantIdLSB(ctxId.tenantId().getId().getLeastSignificantBits()) + .setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits()) + .setEntityType(ctxId.entityId().getEntityType().name()) + .setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits()) + .setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits()) + .build(); + } + + public static CalculatedFieldEntityCtxId fromProto(CalculatedFieldEntityCtxIdProto ctxIdProto) { + TenantId tenantId = TenantId.fromUUID(new UUID(ctxIdProto.getTenantIdMSB(), ctxIdProto.getTenantIdLSB())); + 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); + } + + public static CalculatedFieldStateProto toProto(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state) { + CalculatedFieldStateProto.Builder builder = CalculatedFieldStateProto.newBuilder() + .setId(toProto(stateId)) + .setType(state.getType().name()); + + state.getArguments().forEach((argName, argEntry) -> { + if (argEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { + builder.addSingleValueArguments(toSingleValueArgumentProto(argName, singleValueArgumentEntry)); + } else if (argEntry instanceof TsRollingArgumentEntry rollingArgumentEntry) { + builder.addRollingValueArguments(toRollingArgumentProto(argName, rollingArgumentEntry)); + } + }); + return builder.build(); + } + + public static SingleValueArgumentProto toSingleValueArgumentProto(String argName, SingleValueArgumentEntry entry) { + SingleValueArgumentProto.Builder builder = SingleValueArgumentProto.newBuilder() + .setArgName(argName); + + if (entry.getKvEntryValue() != null) { + builder.setValue(KvProtoUtil.toTsValueProto(entry.getTs(), entry.getKvEntryValue())); + } + + Optional.ofNullable(entry.getVersion()).ifPresent(builder::setVersion); + + return builder.build(); + } + + public static TsRollingArgumentProto toRollingArgumentProto(String argName, TsRollingArgumentEntry entry) { + TsRollingArgumentProto.Builder builder = TsRollingArgumentProto.newBuilder() + .setKey(argName) + .setLimit(entry.getLimit()) + .setTimeWindow(entry.getTimeWindow()); + + entry.getTsRecords().forEach((ts, value) -> builder.addTsValue(TsDoubleValProto.newBuilder().setTs(ts).setValue(value).build())); + + return builder.build(); + } + + public static CalculatedFieldState fromProto(CalculatedFieldStateProto proto) { + if (StringUtils.isEmpty(proto.getType())) { + return null; + } + + CalculatedFieldType type = CalculatedFieldType.valueOf(proto.getType()); + + CalculatedFieldState state = switch (type) { + case SIMPLE -> new SimpleCalculatedFieldState(); + case SCRIPT -> new ScriptCalculatedFieldState(); + }; + + proto.getSingleValueArgumentsList().forEach(argProto -> + state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto))); + + if (CalculatedFieldType.SCRIPT.equals(type)) { + proto.getRollingValueArgumentsList().forEach(argProto -> + state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto))); + } + + return state; + } + + public static SingleValueArgumentEntry fromSingleValueArgumentProto(SingleValueArgumentProto proto) { + if (!proto.hasValue()) { + return new SingleValueArgumentEntry(); + } + TsValueProto tsValueProto = proto.getValue(); + return new SingleValueArgumentEntry( + tsValueProto.getTs(), + (BasicKvEntry) KvProtoUtil.fromTsValueProto(proto.getArgName(), tsValueProto), + proto.getVersion() + ); + } + + public static TsRollingArgumentEntry fromRollingArgumentProto(TsRollingArgumentProto proto) { + TreeMap tsRecords = new TreeMap<>(); + proto.getTsValueList().forEach(tsValueProto -> tsRecords.put(tsValueProto.getTs(), tsValueProto.getValue())); + return new TsRollingArgumentEntry(tsRecords, proto.getLimit(), proto.getTimeWindow()); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/utils/DebugModeRateLimitsConfig.java b/application/src/main/java/org/thingsboard/server/utils/DebugModeRateLimitsConfig.java new file mode 100644 index 0000000000..edae22d1e3 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/DebugModeRateLimitsConfig.java @@ -0,0 +1,36 @@ +/** + * 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.server.utils; + +import lombok.Data; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +@Component +@Data +public class DebugModeRateLimitsConfig { + + @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.enabled:true}") + private boolean ruleChainDebugPerTenantLimitsEnabled; + @Value("${actors.rule.chain.debug_mode_rate_limits_per_tenant.configuration:50000:3600}") + private String ruleChainDebugPerTenantLimitsConfiguration; + + @Value("${actors.calculated_fields.debug_mode_rate_limits_per_tenant.enabled:true}") + private boolean calculatedFieldDebugPerTenantLimitsEnabled; + @Value("${actors.calculated_fields.debug_mode_rate_limits_per_tenant.configuration:50000:3600}") + private String calculatedFieldDebugPerTenantLimitsConfiguration; + +} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index e2570e9137..6cbca3844e 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -509,9 +509,9 @@ actors: calculated_fields: debug_mode_rate_limits_per_tenant: # Enable/Disable the rate limit of persisted debug events for all calculated fields per tenant - enabled: "${ACTORS_CF_DEBUG_MODE_RATE_LIMITS_PER_TENANT_ENABLED:true}" + enabled: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_ENABLED:true}" # The value of DEBUG mode rate limit. By default, no more than 50 thousand events per hour - configuration: "${ACTORS_CF_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}" + configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}" debug: settings: 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 91aad7b1cd..f544f463fa 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -40,8 +40,10 @@ import org.thingsboard.server.controller.CalculatedFieldControllerTest; import org.thingsboard.server.dao.service.DaoSqlTest; import java.util.Map; +import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; @DaoSqlTest public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTest { @@ -84,20 +86,22 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // create CF -> perform initial calculation CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - Thread.sleep(300); - - ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("77.0"); + await().atMost(5, 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}")); - Thread.sleep(300); - - fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + await().atMost(5, 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(); @@ -106,32 +110,35 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes savedOutput.setName("temperatureF"); savedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); - Thread.sleep(300); - - ArrayNode temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); - assertThat(temperatureF).isNotNull(); - assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("86.0"); + await().atMost(5, 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); - Thread.sleep(300); - - temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); - assertThat(temperatureF).isNotNull(); - assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("104.0"); + await().atMost(5, 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); - Thread.sleep(300); - - temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); - assertThat(temperatureF).isNotNull(); - assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("104.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + ArrayNode temperatureF = getServerAttributes(testDevice.getId(), "temperatureF"); + assertThat(temperatureF).isNotNull(); + assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("104.0"); + }); } @Test @@ -164,20 +171,22 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // create CF -> state is not ready -> no calculation performed CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - Thread.sleep(300); - - ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); + await().atMost(5, 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}")); - Thread.sleep(300); - - fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + }); } @Test @@ -211,20 +220,22 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // create CF -> perform initial calculation with default value CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - Thread.sleep(300); - - ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("53.6"); + await().atMost(5, 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}")); - Thread.sleep(300); - - fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + }); } @Test @@ -275,93 +286,100 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // create CF and perform initial calculation doPost("/api/calculatedField", calculatedField, CalculatedField.class); - Thread.sleep(300); - - // result of asset 1 - ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("51.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("51.0"); - // result of asset 2 - ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("52.0"); + // result of asset 2 + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + 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}")); - Thread.sleep(300); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("36.0"); - // result of asset 1 - z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("36.0"); - - // result of asset 2 - z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("37.0"); + // result of asset 2 + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + 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}")); - Thread.sleep(300); - - // result of asset 1 - z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("40.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("40.0"); - // result of asset 2 (no changes) - z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("37.0"); + // result of asset 2 (no changes) + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + 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}")); - Thread.sleep(300); - - // result of asset 1 (no changes) - z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("40.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 (no changes) + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("40.0"); - // result of asset 2 - z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("30.0"); + // result of asset 2 + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + 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}")); - Thread.sleep(300); - - // result of asset 3 - ArrayNode z3 = getServerAttributes(asset3.getId(), "z"); - assertThat(z3).isNotNull(); - assertThat(z3.get(0).get("value").asText()).isEqualTo("38.0"); + Asset finalAsset3 = asset3; + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 3 + ArrayNode z3 = getServerAttributes(finalAsset3.getId(), "z"); + assertThat(z3).isNotNull(); + 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}")); - Thread.sleep(300); - - // result of asset 1 - z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("35.0"); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("35.0"); - // result of asset 2 - z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("25.0"); + // result of asset 2 + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + assertThat(z2.get(0).get("value").asText()).isEqualTo("25.0"); - // result of asset 3 - z3 = getServerAttributes(asset3.getId(), "z"); - assertThat(z3).isNotNull(); - assertThat(z3.get(0).get("value").asText()).isEqualTo("33.0"); + // result of asset 3 + ArrayNode z3 = getServerAttributes(finalAsset3.getId(), "z"); + assertThat(z3).isNotNull(); + assertThat(z3.get(0).get("value").asText()).isEqualTo("33.0"); + }); // update profile for asset 3 -> delete state for asset 3 AssetProfile newAssetProfile = doPost("/api/assetProfile", createAssetProfile("New Asset Profile"), AssetProfile.class); @@ -371,22 +389,24 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // 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}")); - Thread.sleep(300); - - // result of asset 1 - z1 = getServerAttributes(asset1.getId(), "z"); - assertThat(z1).isNotNull(); - assertThat(z1.get(0).get("value").asText()).isEqualTo("30.0"); - - // result of asset 2 - z2 = getServerAttributes(asset2.getId(), "z"); - assertThat(z2).isNotNull(); - assertThat(z2.get(0).get("value").asText()).isEqualTo("20.0"); - - // no changes for asset 3 - z3 = getServerAttributes(asset3.getId(), "z"); - assertThat(z3).isNotNull(); - assertThat(z3.get(0).get("value").asText()).isEqualTo("33.0"); + Asset updatedAsset3 = asset3; + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + // result of asset 1 + ArrayNode z1 = getServerAttributes(asset1.getId(), "z"); + assertThat(z1).isNotNull(); + assertThat(z1.get(0).get("value").asText()).isEqualTo("30.0"); + + // result of asset 2 + ArrayNode z2 = getServerAttributes(asset2.getId(), "z"); + assertThat(z2).isNotNull(); + assertThat(z2.get(0).get("value").asText()).isEqualTo("20.0"); + + // no changes for asset 3 + ArrayNode z3 = getServerAttributes(updatedAsset3.getId(), "z"); + assertThat(z3).isNotNull(); + assertThat(z3.get(0).get("value").asText()).isEqualTo("33.0"); + }); } @Test @@ -421,20 +441,22 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // create CF -> ctx is not initialized -> no calculation perform CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - Thread.sleep(300); - - ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); + await().atMost(5, 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}")); - Thread.sleep(300); - - fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").isNull()).isTrue(); + }); } private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception { 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 f7fc3892b1..bc75bd74ff 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 @@ -20,6 +20,7 @@ import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.mock.mockito.MockBean; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.tbel.DefaultTbelInvokeService; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.common.data.AttributeScope; @@ -35,7 +36,7 @@ import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedField import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.BasicKvEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.service.cf.CalculatedFieldResult; @@ -57,7 +58,7 @@ public class ScriptCalculatedFieldStateTest { private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("5512071d-5abc-411d-a907-4cdb6539c2eb")); private final AssetId ASSET_ID = new AssetId(UUID.fromString("5bc010ae-bcfd-46c8-98b9-8ee8c8955a76")); - private final SingleValueArgumentEntry assetHumidityArgEntry = new SingleValueArgumentEntry(System.currentTimeMillis() - 10, new LongDataEntry("assetHumidity", 43L), 122L); + private final SingleValueArgumentEntry assetHumidityArgEntry = new SingleValueArgumentEntry(System.currentTimeMillis() - 10, new DoubleDataEntry("assetHumidity", 43.0), 122L); private final TsRollingArgumentEntry deviceTemperatureArgEntry = createRollingArgEntry(); private final long ts = System.currentTimeMillis(); @@ -127,52 +128,7 @@ public class ScriptCalculatedFieldStateTest { Output output = getCalculatedFieldConfig().getOutput(); assertThat(result.getType()).isEqualTo(output.getType()); assertThat(result.getScope()).isEqualTo(output.getScope()); - assertThat(result.getResultMap()).isEqualTo(Map.of("averageDeviceTemperature", 13.0, "assetHumidity", 43L)); - } - - @Test - void testPerformCalculationWhenOldTelemetry() throws ExecutionException, InterruptedException { - TsRollingArgumentEntry argumentEntry = new TsRollingArgumentEntry(); - - TreeMap values = new TreeMap<>(); - values.put(ts - 40000, 4.0);// will not be used for calculation - values.put(ts - 45000, 2.0);// will not be used for calculation - values.put(ts - 20, 0.0); - - argumentEntry.setTsRecords(values); - - state.arguments = new HashMap<>(Map.of("deviceTemperature", argumentEntry, "assetHumidity", assetHumidityArgEntry)); - - CalculatedFieldResult result = state.performCalculation(ctx).get(); - - assertThat(result).isNotNull(); - Output output = getCalculatedFieldConfig().getOutput(); - assertThat(result.getType()).isEqualTo(output.getType()); - assertThat(result.getScope()).isEqualTo(output.getScope()); - assertThat(result.getResultMap()).isEqualTo(Map.of("averageDeviceTemperature", 0.0, "assetHumidity", 43L)); - } - - @Test - void testPerformCalculationWhenArgumentsMoreThanLimit() throws ExecutionException, InterruptedException { - TsRollingArgumentEntry argumentEntry = new TsRollingArgumentEntry(); - TreeMap values = new TreeMap<>(); - values.put(ts - 20, 1000.0);// will not be used - values.put(ts - 18, 0.0); - values.put(ts - 16, 0.0); - values.put(ts - 14, 0.0); - values.put(ts - 12, 0.0); - values.put(ts - 10, 0.0); - argumentEntry.setTsRecords(values); - - state.arguments = new HashMap<>(Map.of("deviceTemperature", argumentEntry, "assetHumidity", assetHumidityArgEntry)); - - CalculatedFieldResult result = state.performCalculation(ctx).get(); - - assertThat(result).isNotNull(); - Output output = getCalculatedFieldConfig().getOutput(); - assertThat(result.getType()).isEqualTo(output.getType()); - assertThat(result.getScope()).isEqualTo(output.getScope()); - assertThat(result.getResultMap()).isEqualTo(Map.of("averageDeviceTemperature", 0.0, "assetHumidity", 43L)); + assertThat(result.getResult()).isEqualTo(JacksonUtil.valueToTree(Map.of("maxDeviceTemperature", 17.0, "assetHumidity", 43.0))); } @Test @@ -189,7 +145,7 @@ public class ScriptCalculatedFieldStateTest { @Test void testIsReadyWhenEmptyEntryPresents() { - state.arguments = new HashMap<>(Map.of("deviceTemperature", TsRollingArgumentEntry.EMPTY, "assetHumidity", assetHumidityArgEntry)); + state.arguments = new HashMap<>(Map.of("deviceTemperature", new TsRollingArgumentEntry(5, 30000L), "assetHumidity", assetHumidityArgEntry)); assertThat(state.isReady()).isFalse(); } @@ -235,7 +191,7 @@ public class ScriptCalculatedFieldStateTest { config.setArguments(Map.of("deviceTemperature", argument1, "assetHumidity", argument2)); - config.setExpression("var result = 0; foreach(element : deviceTemperature.entrySet()) { result += element.getValue(); } var map = {}; map.put(\"averageDeviceTemperature\", result / deviceTemperature.size()); map.put(\"assetHumidity\", assetHumidity); return map;"); + config.setExpression("return {\"maxDeviceTemperature\": deviceTemperature.max(), \"assetHumidity\": assetHumidity.value}"); Output output = new Output(); output.setType(OutputType.ATTRIBUTES); diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java index d803ad2dab..55e65422eb 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java @@ -20,6 +20,7 @@ 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.JacksonUtil; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -117,7 +118,7 @@ public class SimpleCalculatedFieldStateTest { "key2", key2ArgEntry )); - Map newArgs = Map.of("key3", TsRollingArgumentEntry.EMPTY); + Map newArgs = Map.of("key3", new TsRollingArgumentEntry(10, 30000L)); assertThatThrownBy(() -> state.updateState(newArgs)) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Rolling argument entry is not supported for simple calculated fields."); @@ -137,7 +138,7 @@ public class SimpleCalculatedFieldStateTest { Output output = getCalculatedFieldConfig().getOutput(); assertThat(result.getType()).isEqualTo(output.getType()); assertThat(result.getScope()).isEqualTo(output.getScope()); - assertThat(result.getResultMap()).isEqualTo(Map.of("output", 49.0)); + assertThat(result.getResult()).isEqualTo(JacksonUtil.valueToTree(Map.of("output", 49.0))); } @Test @@ -175,7 +176,7 @@ public class SimpleCalculatedFieldStateTest { "key1", key1ArgEntry, "key2", key2ArgEntry )); - state.getArguments().put("key3", SingleValueArgumentEntry.EMPTY); + state.getArguments().put("key3", new SingleValueArgumentEntry()); assertThat(state.isReady()).isFalse(); } diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java index 13651e852d..2514b9f34e 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java @@ -40,7 +40,7 @@ public class SingleValueArgumentEntryTest { @Test void testUpdateEntryWhenRollingEntryPassed() { - assertThatThrownBy(() -> entry.updateEntry(TsRollingArgumentEntry.EMPTY)) + assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) .isInstanceOf(IllegalArgumentException.class) .hasMessage("Unsupported argument entry type for single value argument entry: " + ArgumentEntryType.TS_ROLLING); } 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 ad098178f8..75aca61825 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 @@ -17,7 +17,6 @@ package org.thingsboard.server.service.cf.ctx.state; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import org.thingsboard.server.common.data.kv.BasicKvEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; @@ -40,7 +39,7 @@ public class TsRollingArgumentEntryTest { values.put(ts - 30, 12.0); values.put(ts - 20, 17.0); - entry = new TsRollingArgumentEntry(values); + entry = new TsRollingArgumentEntry(5, 30000L, values); } @Test @@ -57,18 +56,10 @@ public class TsRollingArgumentEntryTest { assertThat(entry.getTsRecords().get(ts - 10)).isEqualTo(23.0); } - @Test - void testUpdateEntryWhenSingleValueEntryWithTheSameTsPassed() { - SingleValueArgumentEntry newEntry = new SingleValueArgumentEntry(ts - 20, new DoubleDataEntry("key", 23.0), 123L); - - assertThat(entry.updateEntry(newEntry)).isFalse(); - } - @Test void testUpdateEntryWhenRollingEntryPassed() { TsRollingArgumentEntry newEntry = new TsRollingArgumentEntry(); TreeMap values = new TreeMap<>(); - values.put(ts - 20, 16.0); values.put(ts - 10, 7.0); values.put(ts - 5, 1.0); newEntry.setTsRecords(values); @@ -76,11 +67,11 @@ public class TsRollingArgumentEntryTest { assertThat(entry.updateEntry(newEntry)).isTrue(); assertThat(entry.getTsRecords()).hasSize(5); assertThat(entry.getTsRecords()).isEqualTo(Map.of( - ts - 40, new DoubleDataEntry("key", 10.0), - ts - 30, new DoubleDataEntry("key", 12.0), - ts - 20, new DoubleDataEntry("key", 17.0), - ts - 10, new DoubleDataEntry("key", 7.0), - ts - 5, new DoubleDataEntry("key", 1.0) + ts - 40, 10.0, + ts - 30, 12.0, + ts - 20, 17.0, + ts - 10, 7.0, + ts - 5, 1.0 )); } @@ -90,7 +81,44 @@ public class TsRollingArgumentEntryTest { assertThatThrownBy(() -> entry.updateEntry(newEntry)) .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Argument type " + ArgumentEntryType.TS_ROLLING + " only supports numeric values."); + .hasMessage("Time series rolling arguments supports only numeric values."); + } + + @Test + void testUpdateEntryWhenOldTelemetry() { + TsRollingArgumentEntry newEntry = new TsRollingArgumentEntry(); + TreeMap values = new TreeMap<>(); + values.put(ts - 40000, 4.0);// will not be used for calculation + values.put(ts - 45000, 2.0);// will not be used for calculation + values.put(ts - 5, 0.0); + newEntry.setTsRecords(values); + + entry = new TsRollingArgumentEntry(3, 30000L); + assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.getTsRecords()).hasSize(1); + assertThat(entry.getTsRecords()).isEqualTo(Map.of( + ts - 5, 0.0 + )); + } + + @Test + void testPerformCalculationWhenArgumentsMoreThanLimit() { + TsRollingArgumentEntry newEntry = new TsRollingArgumentEntry(); + TreeMap values = new TreeMap<>(); + values.put(ts - 20, 1000.0);// will not be used + values.put(ts - 18, 0.0); + values.put(ts - 16, 0.0); + values.put(ts - 14, 0.0); + newEntry.setTsRecords(values); + + entry = new TsRollingArgumentEntry(3, 30000L); + assertThat(entry.updateEntry(newEntry)).isTrue(); + assertThat(entry.getTsRecords()).hasSize(3); + assertThat(entry.getTsRecords()).isEqualTo(Map.of( + ts - 18, 0.0, + ts - 16, 0.0, + ts - 14, 0.0 + )); } } \ No newline at end of file diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java index abe1932327..53f39977c1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java @@ -34,4 +34,5 @@ public class SystemParams { boolean mobileQrEnabled; int maxDebugModeDurationMinutes; String ruleChainDebugPerTenantLimitsConfiguration; + String calculatedFieldDebugPerTenantLimitsConfiguration; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEvent.java index e0599db358..9b64dd0d40 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEvent.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEvent.java @@ -91,5 +91,4 @@ public class CalculatedFieldDebugEvent extends Event { return eventInfo; } - } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEventFilter.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEventFilter.java index 839583ef6c..3e56666834 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEventFilter.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/CalculatedFieldDebugEventFilter.java @@ -25,8 +25,6 @@ import org.thingsboard.server.common.data.StringUtils; @Schema public class CalculatedFieldDebugEventFilter extends DebugEventFilter { - @Schema(description = "String value representing the calculated field id in the event body", example = "ccbfa2fe-c8f5-45d8-bb37-6b61a6e02833") - protected String calculatedFieldId; @Schema(description = "String value representing the entity id in the event body", example = "57b6bafe-d600-423c-9267-fe31e5218986") protected String entityId; @Schema(description = "String value representing the entity type", allowableValues = "DEVICE") @@ -35,6 +33,13 @@ public class CalculatedFieldDebugEventFilter extends DebugEventFilter { protected String msgId; @Schema(description = "String value representing the message type", example = "POST_TELEMETRY_REQUEST") protected String msgType; + @Schema(description = "String value representing the arguments that were used in the calculation performed", + example = "{\"x\":{\"ts\":1739432016629,\"value\":20},\"y\":{\"ts\":1739429717656,\"value\":12}}") + protected String arguments; + @Schema(description = "String value representing the result of a calculation", + example = "{\"x + y\":54}") + protected String result; + @Override public EventType getEventType() { @@ -43,8 +48,9 @@ public class CalculatedFieldDebugEventFilter extends DebugEventFilter { @Override public boolean isNotEmpty() { - return super.isNotEmpty() || !StringUtils.isEmpty(calculatedFieldId) || !StringUtils.isEmpty(entityId) - || !StringUtils.isEmpty(entityType) || !StringUtils.isEmpty(msgId) || !StringUtils.isEmpty(msgType); + return super.isNotEmpty() || !StringUtils.isEmpty(entityId) || !StringUtils.isEmpty(entityType) + || !StringUtils.isEmpty(msgId) || !StringUtils.isEmpty(msgType) + || !StringUtils.isEmpty(arguments) || !StringUtils.isEmpty(result); } } 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 6532a4fe0d..35068821ee 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 @@ -43,7 +43,8 @@ public enum LimitedApi { TRANSPORT_MESSAGES_PER_GATEWAY("transport messages per gateway", false), TRANSPORT_MESSAGES_PER_GATEWAY_DEVICE("transport messages per gateway device", false), EMAILS("emails sending", true), - WS_SUBSCRIPTIONS("WS subscriptions", false); + WS_SUBSCRIPTIONS("WS subscriptions", false), + CALCULATED_FIELD_DEBUG_EVENTS(DefaultTenantProfileConfiguration::getCalculatedFieldDebugEventsRateLimit, "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 054adbab1b..0139792fa7 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,6 +140,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxDataPointsPerRollingArg; private long maxStateSizeInKBytes; private long maxSingleValueArgumentSizeInKBytes; + private String calculatedFieldDebugEventsRateLimit; @Override public long getProfileThreshold(ApiUsageRecordKey key) { diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index 07406977d2..cbcbba4fc7 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.ResourceSubType; import org.thingsboard.server.common.data.ResourceType; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; @@ -148,9 +149,13 @@ public class ProtoUtils { var builder = ComponentLifecycleMsg.builder() .tenantId(TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()))) .entityId(entityId) - .event(ComponentLifecycleEvent.values()[proto.getEventValue()]) - .name(proto.getName()) - .oldName(proto.getOldName()); + .event(ComponentLifecycleEvent.values()[proto.getEventValue()]); + if (!StringUtils.isEmpty(proto.getName())) { + builder.name(proto.getName()); + } + if (!StringUtils.isEmpty(proto.getOldName())) { + builder.oldName(proto.getOldName()); + } if (proto.getProfileIdMSB() != 0 || proto.getProfileIdLSB() != 0) { var profileType = EntityType.DEVICE.equals(entityId.getEntityType()) ? EntityType.DEVICE_PROFILE : EntityType.ASSET_PROFILE; builder.profileId(EntityIdFactory.getByTypeAndUuid(profileType, new UUID(proto.getProfileIdMSB(), proto.getProfileIdLSB()))); diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 3dbe01727d..5fc7e7c59c 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -814,11 +814,6 @@ message SingleValueArgumentProto { int64 version = 3; } -message CalculatedFieldStateMsgProto { - CalculatedFieldEntityCtxIdProto id = 1; - CalculatedFieldStateProto state = 2; -} - message TsDoubleValProto { int64 ts = 1; double value = 2; @@ -826,13 +821,16 @@ message TsDoubleValProto { message TsRollingArgumentProto { string key = 1; - repeated TsDoubleValProto tsValue = 2; + int32 limit = 2; + int64 timeWindow = 3; + repeated TsDoubleValProto tsValue = 4; } message CalculatedFieldStateProto { - string type = 1; - repeated SingleValueArgumentProto singleValueArguments = 2; - repeated TsRollingArgumentProto rollingValueArguments = 3; + CalculatedFieldEntityCtxIdProto id = 1; + string type = 2; + repeated SingleValueArgumentProto singleValueArguments = 3; + repeated TsRollingArgumentProto rollingValueArguments = 4; } //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level. diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java index 8e36ae4962..bf830a2ad9 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java @@ -23,7 +23,7 @@ import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueRequestTemplate; @@ -160,12 +160,12 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { + public TbQueueConsumer> createCalculatedFieldStateConsumer() { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); } @Override - public TbQueueProducer> createCalculatedFieldStateProducer() { + public TbQueueProducer> createCalculatedFieldStateProducer() { return new InMemoryTbQueueProducer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index 99a50b2f5a..91ba63427b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -25,7 +25,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.js.JsInvokeProtos; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -548,23 +548,23 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { - return TbKafkaConsumerTemplate.>builder() + public TbQueueConsumer> createCalculatedFieldStateConsumer() { + return TbKafkaConsumerTemplate.>builder() .settings(kafkaSettings) .topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())) .readFromBeginning(true) .stopWhenRead(true) .clientId("monolith-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) .groupId(topicService.buildTopicName("monolith-calculated-field-state-consumer")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), CalculatedFieldStateMsgProto.parseFrom(msg.getData()), msg.getHeaders())) + .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), CalculatedFieldStateProto.parseFrom(msg.getData()), msg.getHeaders())) .admin(cfStateAdmin) .statsService(consumerStatsService) .build(); } @Override - public TbQueueProducer> createCalculatedFieldStateProducer() { - return TbKafkaProducerTemplate.>builder() + public TbQueueProducer> createCalculatedFieldStateProducer() { + return TbKafkaProducerTemplate.>builder() .settings(kafkaSettings) .clientId("monolith-calculated-field-state-" + serviceInfoProvider.getServiceId()) .defaultTopic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index 4fbf5c3e45..d4076ba67d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -23,7 +23,7 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.js.JsInvokeProtos; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -340,23 +340,23 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { - return TbKafkaConsumerTemplate.>builder() + public TbQueueConsumer> createCalculatedFieldStateConsumer() { + return TbKafkaConsumerTemplate.>builder() .settings(kafkaSettings) .topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())) .readFromBeginning(true) .stopWhenRead(true) .clientId("tb-rule-engine-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) .groupId(topicService.buildTopicName("tb-rule-engine-calculated-field-state-consumer")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), CalculatedFieldStateMsgProto.parseFrom(msg.getData()), msg.getHeaders())) + .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), CalculatedFieldStateProto.parseFrom(msg.getData()), msg.getHeaders())) .admin(cfStateAdmin) .statsService(consumerStatsService) .build(); } @Override - public TbQueueProducer> createCalculatedFieldStateProducer() { - return TbKafkaProducerTemplate.>builder() + public TbQueueProducer> createCalculatedFieldStateProducer() { + return TbKafkaProducerTemplate.>builder() .settings(kafkaSettings) .clientId("tb-rule-engine-to-calculated-field-state-" + serviceInfoProvider.getServiceId()) .defaultTopic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java index 5f0af72955..d3c1d09399 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java @@ -17,7 +17,7 @@ package org.thingsboard.server.queue.provider; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.gen.js.JsInvokeProtos; -import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -126,8 +126,8 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueConsumer> createToCalculatedFieldNotificationsMsgConsumer(); - TbQueueConsumer> createCalculatedFieldStateConsumer(); + TbQueueConsumer> createCalculatedFieldStateConsumer(); - TbQueueProducer> createCalculatedFieldStateProducer(); + TbQueueProducer> createCalculatedFieldStateProducer(); } diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/DefaultTbelInvokeService.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/DefaultTbelInvokeService.java index 177bc517a8..2d19e39c31 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/DefaultTbelInvokeService.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/DefaultTbelInvokeService.java @@ -136,6 +136,7 @@ public class DefaultTbelInvokeService extends AbstractScriptInvokeService implem parserConfig.registerDataType("TbelCfSingleValueArg", TbelCfSingleValueArg.class, TbelCfSingleValueArg::memorySize); parserConfig.registerDataType("TbelCfTsRollingArg", TbelCfTsRollingArg.class, TbelCfTsRollingArg::memorySize); parserConfig.registerDataType("TbelCfTsDoubleVal", TbelCfTsDoubleVal.class, TbelCfTsDoubleVal::memorySize); + parserConfig.registerDataType("TbTimeWindow", TbTimeWindow.class, TbTimeWindow::memorySize); TbUtils.register(parserConfig); executor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "tbel-executor")); try { diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbTimeWindow.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbTimeWindow.java new file mode 100644 index 0000000000..f90a60ff3a --- /dev/null +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbTimeWindow.java @@ -0,0 +1,36 @@ +/** + * 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 lombok.AllArgsConstructor; +import lombok.Data; + +@Data +@AllArgsConstructor +public class TbTimeWindow implements TbelCfObject { + + public static final long OBJ_SIZE = 32L; + + private long startTs; + private long endTs; + private int limit; + + @Override + public long memorySize() { + return OBJ_SIZE; + } + +} 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 ac5bd5cc06..7fe75acd05 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,20 @@ */ package org.thingsboard.script.api.tbel; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; + +@JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + include = JsonTypeInfo.As.PROPERTY, + property = "type" +) +@JsonSubTypes({ + @JsonSubTypes.Type(value = TbelCfSingleValueArg.class, name = "SINGLE_VALUE"), + @JsonSubTypes.Type(value = TbelCfTsRollingArg.class, name = "TS_ROLLING") +}) public interface TbelCfArg extends TbelCfObject { + String getType(); + } diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfSingleValueArg.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfSingleValueArg.java index 7076e9dd5e..d35ac36ab0 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfSingleValueArg.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfSingleValueArg.java @@ -15,6 +15,8 @@ */ package org.thingsboard.script.api.tbel; +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; import lombok.Data; @Data @@ -23,9 +25,23 @@ public class TbelCfSingleValueArg implements TbelCfArg { private final long ts; private final Object value; + @JsonCreator + public TbelCfSingleValueArg( + @JsonProperty("ts") long ts, + @JsonProperty("value") Object value + ) { + this.ts = ts; + this.value = value; + } + @Override public long memorySize() { return 8L; // TODO; } + @Override + public String getType() { + return "SINGLE_VALUE"; + } + } 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 14122bbde8..843abf8d5c 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 @@ -15,7 +15,9 @@ */ package org.thingsboard.script.api.tbel; +import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonProperty; import lombok.Getter; import java.util.Collections; @@ -27,10 +29,23 @@ import static org.thingsboard.script.api.tbel.TbelCfTsDoubleVal.OBJ_SIZE; public class TbelCfTsRollingArg implements TbelCfArg, Iterable { + @Getter + private final TbTimeWindow timeWindow; @Getter private final List values; - public TbelCfTsRollingArg(List values) { + @JsonCreator + public TbelCfTsRollingArg( + @JsonProperty("timeWindow") TbTimeWindow timeWindow, + @JsonProperty("values") List values + ) { + this.timeWindow = timeWindow; + this.values = Collections.unmodifiableList(values); + } + + public TbelCfTsRollingArg(int limit, long timeWindow, List values) { + long ts = System.currentTimeMillis(); + this.timeWindow = new TbTimeWindow(ts - timeWindow, ts, limit); this.values = Collections.unmodifiableList(values); } @@ -70,4 +85,9 @@ public class TbelCfTsRollingArg implements TbelCfArg, Iterable implements BaseEntity { +public class CalculatedFieldEntity extends BaseVersionedEntity implements BaseEntity { @Column(name = CALCULATED_FIELD_TENANT_ID_COLUMN) private UUID tenantId; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/CalculatedFieldDebugEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/CalculatedFieldDebugEventRepository.java index a759babcfb..56a39c41f2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/CalculatedFieldDebugEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/CalculatedFieldDebugEventRepository.java @@ -59,6 +59,8 @@ public interface CalculatedFieldDebugEventRepository extends EventRepository findEventByFilter(UUID tenantId, UUID entityId, CalculatedFieldDebugEventFilter eventFilter, TimePageLink pageLink) { - parseUUID(eventFilter.getCalculatedFieldId(), "Calculated Field Id"); parseUUID(eventFilter.getEntityId(), "Entity Id"); parseUUID(eventFilter.getMsgId(), "Message Id"); return DaoUtil.toPageData( @@ -304,11 +303,13 @@ public class JpaBaseEventDao implements EventDao { pageLink.getStartTime(), pageLink.getEndTime(), eventFilter.getServer(), - UUID.fromString(eventFilter.getCalculatedFieldId()), + entityId, eventFilter.getEntityId(), eventFilter.getEntityType(), eventFilter.getMsgId(), eventFilter.getMsgType(), + eventFilter.getArguments(), + eventFilter.getResult(), eventFilter.isError(), eventFilter.getErrorStr(), DaoUtil.toPageable(pageLink, EventEntity.eventColumnMap))); @@ -389,7 +390,6 @@ public class JpaBaseEventDao implements EventDao { } private void removeEventsByFilter(UUID tenantId, UUID entityId, CalculatedFieldDebugEventFilter eventFilter, Long startTime, Long endTime) { - parseUUID(eventFilter.getCalculatedFieldId(), "Calculated Field Id"); parseUUID(eventFilter.getEntityId(), "Entity Id"); parseUUID(eventFilter.getMsgId(), "Message Id"); calculatedFieldDebugEventRepository.removeEvents( @@ -398,11 +398,13 @@ public class JpaBaseEventDao implements EventDao { startTime, endTime, eventFilter.getServer(), - UUID.fromString(eventFilter.getCalculatedFieldId()), + entityId, eventFilter.getEntityId(), eventFilter.getEntityType(), eventFilter.getMsgId(), eventFilter.getMsgType(), + eventFilter.getArguments(), + eventFilter.getResult(), eventFilter.isError(), eventFilter.getErrorStr()); }