From 958e44f7d667ce697561f48509f41e372fc1c505 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 5 Dec 2025 14:14:17 +0200 Subject: [PATCH 1/8] set ts -1 when default value used --- ...faultCalculatedFieldProcessingService.java | 68 ++++------------- .../ctx/state/BaseCalculatedFieldState.java | 5 ++ .../ctx/state/SimpleCalculatedFieldState.java | 5 +- .../ctx/state/SingleValueArgumentEntry.java | 6 +- .../utils/CalculatedFieldArgumentUtils.java | 75 +++++++++++++++++++ .../cf/CalculatedFieldIntegrationTest.java | 17 ++++- .../script/api/tbel/TbelCfCtx.java | 2 +- 7 files changed, 112 insertions(+), 66 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java 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 f2a6916751..9ff185bb52 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 @@ -23,7 +23,6 @@ import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.math.NumberUtils; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; @@ -31,21 +30,14 @@ import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.OutputType; 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.kv.Aggregation; -import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.BooleanDataEntry; -import org.thingsboard.server.common.data.kv.DoubleDataEntry; -import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; -import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; @@ -70,9 +62,6 @@ 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; -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.ArrayList; @@ -80,12 +69,16 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Map.Entry; -import java.util.Optional; import java.util.UUID; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.SCOPE; +import static org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry.DEFAULT_TS; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultTsKvEntry; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; @TbRuleEngineComponent @@ -244,30 +237,17 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { return switch (argument.getRefEntityKey().getType()) { case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); - case ATTRIBUTE -> transformSingleValueArgument( - Futures.transform( - attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), - result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), - calculatedFieldCallbackExecutor) - ); - case TS_LATEST -> transformSingleValueArgument( - Futures.transform( - timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), - result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))), - calculatedFieldCallbackExecutor)); + case ATTRIBUTE -> Futures.transform( + attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), + result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultAttributeEntry(argument, DEFAULT_TS))), + calculatedFieldCallbackExecutor); + case TS_LATEST -> Futures.transform( + timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), + result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultTsKvEntry(argument, DEFAULT_TS))), + calculatedFieldCallbackExecutor); }; } - private ListenableFuture transformSingleValueArgument(ListenableFuture> kvEntryFuture) { - return Futures.transform(kvEntryFuture, kvEntry -> { - if (kvEntry.isPresent() && kvEntry.get().getValue() != null) { - return ArgumentEntry.createSingleValueArgument(kvEntry.get()); - } else { - return new SingleValueArgumentEntry(); - } - }, calculatedFieldCallbackExecutor); - } - private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) { long currentTime = System.currentTimeMillis(); long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); @@ -282,28 +262,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? new TsRollingArgumentEntry(limit, timeWindow) : ArgumentEntry.createTsRollingArgument(tsRolling, limit, timeWindow), calculatedFieldCallbackExecutor); } - private KvEntry createDefaultKvEntry(Argument argument) { - String key = argument.getRefEntityKey().getKey(); - String defaultValue = argument.getDefaultValue(); - if (StringUtils.isBlank(defaultValue)) { - return new StringDataEntry(key, null); - } - if (NumberUtils.isParsable(defaultValue)) { - return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); - } - if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { - return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); - } - return new StringDataEntry(key, defaultValue); - } - - private CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) { - return switch (ctx.getCfType()) { - case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames()); - case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames()); - }; - } - private static class TbCallbackWrapper implements TbQueueCallback { private final TbCallback callback; 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 94a18256a3..856081d7c5 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 @@ -48,6 +48,11 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { this(new ArrayList<>(), new HashMap<>(), false, DEFAULT_LAST_UPDATE_TS); } + + public long getLatestTimestamp() { + return latestTimestamp == DEFAULT_LAST_UPDATE_TS ? System.currentTimeMillis() : latestTimestamp; + } + @Override public boolean updateState(CalculatedFieldCtx ctx, Map argumentValues) { if (arguments == null) { 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 80f5964582..e80939a952 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 @@ -99,10 +99,9 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { valuesNode.set(outputName, JacksonUtil.valueToTree(result)); } - long latestTs = getLatestTimestamp(); - if (useLatestTs && latestTs != DEFAULT_LAST_UPDATE_TS) { + if (useLatestTs) { ObjectNode resultNode = JacksonUtil.newObjectNode(); - resultNode.put("ts", latestTs); + resultNode.put("ts", getLatestTimestamp()); resultNode.set("values", valuesNode); return resultNode; } else { 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 0997fd6cbb..b26e53ef3c 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 @@ -19,7 +19,6 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.core.type.TypeReference; import lombok.AllArgsConstructor; import lombok.Data; -import lombok.NoArgsConstructor; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; @@ -32,13 +31,12 @@ import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.gen.transport.TransportProtos.AttributeValueProto; import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; -import static org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState.DEFAULT_LAST_UPDATE_TS; - @Data @AllArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { public static final Long DEFAULT_VERSION = -1L; + public static final Long DEFAULT_TS = -1L; private long ts; private BasicKvEntry kvEntryValue; @@ -47,7 +45,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { private boolean forceResetPrevious; public SingleValueArgumentEntry() { - this.ts = DEFAULT_LAST_UPDATE_TS; + this.ts = DEFAULT_TS; this.version = DEFAULT_VERSION; } diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java new file mode 100644 index 0000000000..1296df6e87 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -0,0 +1,75 @@ +/** + * Copyright © 2016-2025 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.NonNull; +import org.apache.commons.lang3.math.NumberUtils; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +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; +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 static org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry.DEFAULT_VERSION; + +public class CalculatedFieldArgumentUtils { + + public static ArgumentEntry transformSingleValueArgument(@NonNull KvEntry kvEntry) { + return kvEntry.getValue() != null ? ArgumentEntry.createSingleValueArgument(kvEntry) : new SingleValueArgumentEntry(); + } + + public static TsKvEntry createDefaultTsKvEntry(Argument argument, long ts) { + return new BasicTsKvEntry(ts, createDefaultKvEntry(argument), DEFAULT_VERSION); + } + + public static AttributeKvEntry createDefaultAttributeEntry(Argument argument, long ts) { + return new BaseAttributeKvEntry(createDefaultKvEntry(argument), ts, DEFAULT_VERSION); + } + + private static KvEntry createDefaultKvEntry(Argument argument) { + String key = argument.getRefEntityKey().getKey(); + String defaultValue = argument.getDefaultValue(); + if (StringUtils.isBlank(defaultValue)) { + return new StringDataEntry(key, null); + } + if (NumberUtils.isParsable(defaultValue)) { + return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); + } + if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { + return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); + } + return new StringDataEntry(key, defaultValue); + } + + public static CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) { + return switch (ctx.getCfType()) { + case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames()); + case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames()); + }; + } + +} 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 b500f95d45..03052382f0 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -570,8 +570,6 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes @Test public void testScriptCalculatedFieldWhenUsedLatestTsInScript() throws Exception { Device testDevice = createDevice("Test device", "1234567890"); - long ts = System.currentTimeMillis() - 300000L; - doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); CalculatedField calculatedField = new CalculatedField(); calculatedField.setEntityId(testDevice.getId()); @@ -585,6 +583,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes Argument argument = new Argument(); ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); argument.setRefEntityKey(refEntityKey); + argument.setDefaultValue("20"); config.setArguments(Map.of("T", argument)); config.setExpression("return {\"ts\": ctx.latestTs, \"values\": {\"fahrenheitTemp\": (T * 1.8) + 32}};"); @@ -596,7 +595,19 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + await().alias("create CF -> perform initial calculation with default value").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isNotEqualTo("-1"); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("68.0"); + }); + + long ts = System.currentTimeMillis() - 10L; + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); + + await().alias("update telemetry -> perform calculation").atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java index c6023154ea..2fe861ba81 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java @@ -29,7 +29,7 @@ public class TbelCfCtx implements TbelCfObject { public TbelCfCtx(Map args, long latestTs) { this.args = Collections.unmodifiableMap(args); - this.latestTs = latestTs != -1 ? latestTs : System.currentTimeMillis(); + this.latestTs = latestTs; } @Override From 4e2b4fc9216b554343704fa5b4ddb9d5113b5710 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 8 Dec 2025 10:08:18 +0200 Subject: [PATCH 2/8] added getLatestTs method and added tests --- ...faultCalculatedFieldProcessingService.java | 5 +- .../service/cf/ctx/state/ArgumentEntry.java | 2 + .../ctx/state/BaseCalculatedFieldState.java | 22 ++--- .../ctx/state/SimpleCalculatedFieldState.java | 5 +- .../ctx/state/SingleValueArgumentEntry.java | 17 +++- .../cf/ctx/state/TsRollingArgumentEntry.java | 8 ++ .../cf/CalculatedFieldIntegrationTest.java | 98 ++++++++++++++++--- .../script/api/tbel/TbelCfCtx.java | 2 +- 8 files changed, 124 insertions(+), 35 deletions(-) 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 9ff185bb52..b6e1193cfb 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 @@ -74,7 +74,6 @@ import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.SCOPE; -import static org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry.DEFAULT_TS; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultTsKvEntry; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; @@ -239,11 +238,11 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); case ATTRIBUTE -> Futures.transform( attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), - result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultAttributeEntry(argument, DEFAULT_TS))), + result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultAttributeEntry(argument, System.currentTimeMillis()))), calculatedFieldCallbackExecutor); case TS_LATEST -> Futures.transform( timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), - result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultTsKvEntry(argument, DEFAULT_TS))), + result -> transformSingleValueArgument(result.orElseGet(() -> createDefaultTsKvEntry(argument, System.currentTimeMillis()))), calculatedFieldCallbackExecutor); }; } 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 83e10b8194..3df43d8c2b 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 @@ -40,6 +40,8 @@ public interface ArgumentEntry { Object getValue(); + long getLatestTs(); + boolean updateEntry(ArgumentEntry entry); boolean isEmpty(); 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 856081d7c5..b6d4ad0fac 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 @@ -48,11 +48,6 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { this(new ArrayList<>(), new HashMap<>(), false, DEFAULT_LAST_UPDATE_TS); } - - public long getLatestTimestamp() { - return latestTimestamp == DEFAULT_LAST_UPDATE_TS ? System.currentTimeMillis() : latestTimestamp; - } - @Override public boolean updateState(CalculatedFieldCtx ctx, Map argumentValues) { if (arguments == null) { @@ -80,7 +75,6 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { if (entryUpdated) { stateUpdated = true; - updateLastUpdateTimestamp(newEntry); } } @@ -116,15 +110,13 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected abstract void validateNewEntry(ArgumentEntry newEntry); - private void updateLastUpdateTimestamp(ArgumentEntry entry) { - long newTs = this.latestTimestamp; - if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - newTs = singleValueArgumentEntry.getTs(); - } else if (entry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().lastEntry(); - newTs = (lastEntry != null) ? lastEntry.getKey() : DEFAULT_LAST_UPDATE_TS; - } - this.latestTimestamp = Math.max(this.latestTimestamp, newTs); + public long getLatestTimestamp() { + long currentLatestTs = arguments.values().stream() + .mapToLong(ArgumentEntry::getLatestTs) + .max() + .orElse(DEFAULT_LAST_UPDATE_TS); + latestTimestamp = Math.max(currentLatestTs, latestTimestamp); + return latestTimestamp; } } 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 e80939a952..80f5964582 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 @@ -99,9 +99,10 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { valuesNode.set(outputName, JacksonUtil.valueToTree(result)); } - if (useLatestTs) { + long latestTs = getLatestTimestamp(); + if (useLatestTs && latestTs != DEFAULT_LAST_UPDATE_TS) { ObjectNode resultNode = JacksonUtil.newObjectNode(); - resultNode.put("ts", getLatestTimestamp()); + resultNode.put("ts", latestTs); resultNode.set("values", valuesNode); return resultNode; } else { 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 b26e53ef3c..186af92f38 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 @@ -31,12 +31,13 @@ import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.gen.transport.TransportProtos.AttributeValueProto; import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; +import static org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState.DEFAULT_LAST_UPDATE_TS; + @Data @AllArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { public static final Long DEFAULT_VERSION = -1L; - public static final Long DEFAULT_TS = -1L; private long ts; private BasicKvEntry kvEntryValue; @@ -45,7 +46,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { private boolean forceResetPrevious; public SingleValueArgumentEntry() { - this.ts = DEFAULT_TS; + this.ts = DEFAULT_LAST_UPDATE_TS; this.version = DEFAULT_VERSION; } @@ -97,6 +98,11 @@ public class SingleValueArgumentEntry implements ArgumentEntry { return isEmpty() ? null : kvEntryValue.getValue(); } + @Override + public long getLatestTs() { + return !isDefaultValue() ? ts : DEFAULT_LAST_UPDATE_TS; + } + @Override public TbelCfArg toTbelCfArg() { Object value = kvEntryValue.getValue(); @@ -118,7 +124,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { @Override public boolean updateEntry(ArgumentEntry entry) { if (entry instanceof SingleValueArgumentEntry singleValueEntry) { - if (singleValueEntry.getTs() <= this.ts) { + if (singleValueEntry.getTs() < this.ts) { return false; } @@ -134,4 +140,9 @@ public class SingleValueArgumentEntry implements ArgumentEntry { } return false; } + + public boolean isDefaultValue() { + return DEFAULT_VERSION.equals(this.version); + } + } 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 b5a680a072..ada46a841e 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 @@ -31,6 +31,8 @@ import java.util.List; import java.util.Map; import java.util.TreeMap; +import static org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState.DEFAULT_LAST_UPDATE_TS; + @Data @NoArgsConstructor @AllArgsConstructor @@ -83,6 +85,12 @@ public class TsRollingArgumentEntry implements ArgumentEntry { return tsRecords; } + @Override + public long getLatestTs() { + var lastEntry = tsRecords.lastEntry(); + return (lastEntry != null) ? lastEntry.getKey() : DEFAULT_LAST_UPDATE_TS; + } + @Override public TbelCfArg toTbelCfArg() { List values = new ArrayList<>(tsRecords.size()); 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 03052382f0..c83e8c5ee3 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -45,6 +45,7 @@ import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @DaoSqlTest public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTest { @@ -571,6 +572,9 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes public void testScriptCalculatedFieldWhenUsedLatestTsInScript() throws Exception { Device testDevice = createDevice("Test device", "1234567890"); + long ts = System.currentTimeMillis() - 300000L; + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); + CalculatedField calculatedField = new CalculatedField(); calculatedField.setEntityId(testDevice.getId()); calculatedField.setType(CalculatedFieldType.SCRIPT); @@ -583,7 +587,6 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes Argument argument = new Argument(); ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); argument.setRefEntityKey(refEntityKey); - argument.setDefaultValue("20"); config.setArguments(Map.of("T", argument)); config.setExpression("return {\"ts\": ctx.latestTs, \"values\": {\"fahrenheitTemp\": (T * 1.8) + 32}};"); @@ -595,25 +598,98 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); - await().alias("create CF -> perform initial calculation with default value").atMost(TIMEOUT, TimeUnit.SECONDS) + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isNotEqualTo("-1"); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("68.0"); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); }); + } - long ts = System.currentTimeMillis() - 10L; - doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); + @Test + public void testSimpleCalculatedFieldWhenUseLatestTsIsTrueAndDefaultArguments() throws Exception { + Device testDevice = createDevice("Test device", "1234567890"); - await().alias("update telemetry -> perform calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(testDevice.getId()); + calculatedField.setType(CalculatedFieldType.SIMPLE); + calculatedField.setName("a + b + c"); + calculatedField.setDebugSettings(DebugSettings.all()); + calculatedField.setConfigurationVersion(1); + + SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); + + Argument argument1 = new Argument(); + ReferencedEntityKey refEntityKey1 = new ReferencedEntityKey("a", ArgumentType.TS_LATEST, null); + argument1.setRefEntityKey(refEntityKey1); + argument1.setDefaultValue("100"); + Argument argument2 = new Argument(); + ReferencedEntityKey refEntityKey2 = new ReferencedEntityKey("b", ArgumentType.TS_LATEST, null); + argument2.setRefEntityKey(refEntityKey2); + argument2.setDefaultValue("200"); + Argument argument3 = new Argument(); + ReferencedEntityKey refEntityKey3 = new ReferencedEntityKey("c", ArgumentType.TS_LATEST, null); + argument3.setRefEntityKey(refEntityKey3); + argument3.setDefaultValue("300"); + config.setArguments(Map.of("a", argument1, "b", argument2, "c", argument3)); + config.setExpression("a + b + c"); + + Output output = new Output(); + output.setName("d"); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(0); + config.setOutput(output); + + config.setUseLatestTs(true); + + calculatedField.setConfiguration(config); + + CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); + + await().alias("create CF -> perform initial calculation with default arguments").atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); - assertThat(fahrenheitTemp).isNotNull(); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); - assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + ObjectNode d = getLatestTelemetry(testDevice.getId(), "d"); + assertThat(d).isNotNull(); + assertThat(d.get("d").get(0).get("value").asText()).isEqualTo("600"); + }); + + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"a\":10}")); + + await().alias("update telemetry -> save result with ts of 'a' argument").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode keys = getLatestTelemetry(testDevice.getId(), "d", "a"); + assertThat(keys).isNotNull(); + String aTs = keys.get("a").get(0).get("ts").asText(); + assertThat(keys.get("d").get(0).get("ts").asText()).isEqualTo(aTs); + assertThat(keys.get("d").get(0).get("value").asText()).isEqualTo("510"); + }); + + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"b\":20}")); + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"c\":30}")); + + await().alias("update telemetry -> save result with latest ts of updated arguments").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode keys = getLatestTelemetry(testDevice.getId(), "d"); + assertThat(keys).isNotNull(); + assertThat(keys.get("d").get(0).get("value").asText()).isEqualTo("60"); + }); + + String latestTs = getLatestTelemetry(testDevice.getId(), "d").get("d").get(0).get("ts").asText(); + + doDelete("/api/plugins/telemetry/DEVICE/" + testDevice.getId() + "/timeseries/delete?keys=b&deleteAllDataForKeys=true").andExpect(status().isOk()); + + await().alias("delete telemetry -> save result with previous latest ts and default argument").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode keys = getLatestTelemetry(testDevice.getId(), "d"); + assertThat(keys).isNotNull(); + assertThat(keys.get("d").get(0).get("ts").asText()).isEqualTo(latestTs); + assertThat(keys.get("d").get(0).get("value").asText()).isEqualTo("240"); }); } diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java index 2fe861ba81..c6023154ea 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java @@ -29,7 +29,7 @@ public class TbelCfCtx implements TbelCfObject { public TbelCfCtx(Map args, long latestTs) { this.args = Collections.unmodifiableMap(args); - this.latestTs = latestTs; + this.latestTs = latestTs != -1 ? latestTs : System.currentTimeMillis(); } @Override From 0577d66dcad395f3c82a1e8a547e54df1f5a5b9b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 8 Dec 2025 12:37:20 +0200 Subject: [PATCH 3/8] return the most recent ts of default args if all are default --- .../service/cf/ctx/state/ArgumentEntry.java | 2 -- .../ctx/state/BaseCalculatedFieldState.java | 32 +++++++++++++------ .../ctx/state/SingleValueArgumentEntry.java | 5 --- .../cf/ctx/state/TsRollingArgumentEntry.java | 1 - 4 files changed, 23 insertions(+), 17 deletions(-) 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 3df43d8c2b..83e10b8194 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 @@ -40,8 +40,6 @@ public interface ArgumentEntry { Object getValue(); - long getLatestTs(); - boolean updateEntry(ArgumentEntry entry); boolean isEmpty(); 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 b6d4ad0fac..14370aa68c 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 @@ -37,15 +37,13 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected Map arguments; protected boolean sizeExceedsLimit; - protected long latestTimestamp = DEFAULT_LAST_UPDATE_TS; - public BaseCalculatedFieldState(List requiredArguments) { this.requiredArguments = requiredArguments; this.arguments = new HashMap<>(); } public BaseCalculatedFieldState() { - this(new ArrayList<>(), new HashMap<>(), false, DEFAULT_LAST_UPDATE_TS); + this(new ArrayList<>(), new HashMap<>(), false); } @Override @@ -111,12 +109,28 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected abstract void validateNewEntry(ArgumentEntry newEntry); public long getLatestTimestamp() { - long currentLatestTs = arguments.values().stream() - .mapToLong(ArgumentEntry::getLatestTs) - .max() - .orElse(DEFAULT_LAST_UPDATE_TS); - latestTimestamp = Math.max(currentLatestTs, latestTimestamp); - return latestTimestamp; + long latestTs = DEFAULT_LAST_UPDATE_TS; + + boolean allDefault = arguments.values().stream().allMatch(entry -> { + if (entry instanceof SingleValueArgumentEntry single) { + return single.isDefaultValue(); + } + return false; + }); + + for (ArgumentEntry entry : arguments.values()) { + if (entry instanceof SingleValueArgumentEntry single) { + if (allDefault) { + latestTs = Math.max(latestTs, single.getTs()); + } else if (!single.isDefaultValue()) { + latestTs = Math.max(latestTs, single.getTs()); + } + } else if (entry instanceof TsRollingArgumentEntry rolling) { + latestTs = Math.max(latestTs, rolling.getLatestTs()); + } + } + + return latestTs; } } 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 186af92f38..2f9a7de940 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 @@ -98,11 +98,6 @@ public class SingleValueArgumentEntry implements ArgumentEntry { return isEmpty() ? null : kvEntryValue.getValue(); } - @Override - public long getLatestTs() { - return !isDefaultValue() ? ts : DEFAULT_LAST_UPDATE_TS; - } - @Override public TbelCfArg toTbelCfArg() { Object value = kvEntryValue.getValue(); 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 ada46a841e..e01d8b7369 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 @@ -85,7 +85,6 @@ public class TsRollingArgumentEntry implements ArgumentEntry { return tsRecords; } - @Override public long getLatestTs() { var lastEntry = tsRecords.lastEntry(); return (lastEntry != null) ? lastEntry.getKey() : DEFAULT_LAST_UPDATE_TS; From 09eb511599225583100fe2eb2ab7c1f79bc98f48 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 8 Dec 2025 16:51:40 +0200 Subject: [PATCH 4/8] Fix key_dictionary race causing cached keyId 0 --- .../dictionary/KeyDictionaryCompositeKey.java | 2 +- .../sqlts/dictionary/KeyDictionaryEntry.java | 5 +- .../sqlts/dictionary/JpaKeyDictionaryDao.java | 59 ++++------ .../dictionary/KeyDictionaryRepository.java | 4 + .../dictionary/KeyDictionaryDaoTest.java | 111 ++++++++++++++++++ 5 files changed, 141 insertions(+), 40 deletions(-) create mode 100644 dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java index 4f3285b9bf..00e49ea703 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java @@ -25,7 +25,7 @@ import java.io.Serializable; @Data @NoArgsConstructor @AllArgsConstructor -public class KeyDictionaryCompositeKey implements Serializable{ +public class KeyDictionaryCompositeKey implements Serializable { @Transient private static final long serialVersionUID = -4089175869616037523L; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java index a95c7a2bc6..d98105c0bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java @@ -36,8 +36,7 @@ public final class KeyDictionaryEntry { @Column(name = KEY_COLUMN) private String key; - @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int") - @Generated - private int keyId; + @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int", insertable = false, updatable = false) + private Integer keyId; } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index c14c069f23..1ef2c2d6aa 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -17,8 +17,6 @@ package org.thingsboard.server.dao.sqlts.dictionary; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.hibernate.exception.ConstraintViolationException; -import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; @@ -48,43 +46,32 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService @Transactional(propagation = Propagation.NOT_SUPPORTED) @Override public Integer getOrSaveKeyId(String strKey) { - Integer keyId = keyDictionaryMap.get(strKey); - if (keyId == null) { - Optional tsKvDictionaryOptional; - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - creationLock.lock(); - try { - keyId = keyDictionaryMap.get(strKey); - if (keyId != null) { - return keyId; - } - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - KeyDictionaryEntry keyDictionaryEntry = new KeyDictionaryEntry(); - keyDictionaryEntry.setKey(strKey); - try { - KeyDictionaryEntry saved = keyDictionaryRepository.save(keyDictionaryEntry); - keyDictionaryMap.put(saved.getKey(), saved.getKeyId()); - keyId = saved.getKeyId(); - } catch (DataIntegrityViolationException | ConstraintViolationException e) { - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - KeyDictionaryEntry dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get KeyDictionaryEntry entity from DB!")); - keyDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); - keyId = dictionary.getKeyId(); - } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - } - } finally { - creationLock.unlock(); + Integer cached = keyDictionaryMap.get(strKey); + if (cached != null) { + return cached; + } + creationLock.lock(); + try { + Integer keyId = keyDictionaryMap.get(strKey); + if (keyId != null) { + return keyId; + } + keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); + if (keyId == null || keyId == 0) { + log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); + KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(strKey); + Optional entryOpt = keyDictionaryRepository.findById(id); + if (entryOpt.isEmpty() || + entryOpt.get().getKeyId() == null || + entryOpt.get().getKeyId() == 0) { + throw new IllegalStateException("Failed to resolve keyId for string key: " + strKey + " after fallback. keyId: " + keyId); } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - keyDictionaryMap.put(strKey, keyId); } + keyDictionaryMap.put(strKey, keyId); + return keyId; + } finally { + creationLock.unlock(); } - return keyId; } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java index d264cd9966..a1141b62ff 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java @@ -19,6 +19,7 @@ import org.springframework.data.domain.Page; import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; @@ -31,4 +32,7 @@ public interface KeyDictionaryRepository extends JpaRepository findAll(Pageable pageable); + @Query(value = "INSERT INTO key_dictionary (key) VALUES (:key) ON CONFLICT (key) DO UPDATE SET key = EXCLUDED.key RETURNING key_id", nativeQuery = true) + Integer upsertAndGetKeyId(@Param("key") String key); + } \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java new file mode 100644 index 0000000000..11c14c99f4 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java @@ -0,0 +1,111 @@ +/** + * Copyright © 2016-2025 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.dao.sqlts.dictionary; + +import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.dao.dictionary.KeyDictionaryDao; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; +import org.thingsboard.server.dao.service.AbstractServiceTest; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.Arrays; +import java.util.Optional; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; + +@DaoSqlTest +public class KeyDictionaryDaoTest extends AbstractServiceTest { + + @Autowired + private KeyDictionaryDao keyDictionaryDao; + + @Autowired + private KeyDictionaryRepository keyDictionaryRepository; + + private static final String KEY = "testKeyDictionaryDaoTestKey"; + + @Test + public void testGetOrSaveKeyId_concurrent() throws Exception { + int threads = 8; + ExecutorService executor = Executors.newFixedThreadPool(threads); + + CountDownLatch allReady = new CountDownLatch(threads); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch allDone = new CountDownLatch(threads); + + Integer[] keyIds = new Integer[threads]; + + try { + for (int i = 0; i < threads; i++) { + final int idx = i; + executor.submit(() -> { + allReady.countDown(); + try { + // wait until all threads are ready + start.await(); + // concurrent call + Integer id = keyDictionaryDao.getOrSaveKeyId(KEY); + keyIds[idx] = id; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + allDone.countDown(); + } + }); + } + + // ensure all threads are queued + allReady.await(5, TimeUnit.SECONDS); + // fire the start gun + start.countDown(); + // wait for all to finish + allDone.await(10, TimeUnit.SECONDS); + } finally { + executor.shutdownNow(); + } + + // basic sanity + for (int i = 0; i < threads; i++) { + assertThat(keyIds[i]) + .as("keyId[%s]", i) + .isNotNull() + .isGreaterThan(0); + } + + // all threads must see the same keyId + int first = keyIds[0]; + assertThat(first).isGreaterThan(0); + assertThat(Arrays.stream(keyIds).distinct().count()) + .as("all threads should get the same keyId") + .isEqualTo(1); + + // DB must have exactly one row for this key and the same id + KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(KEY); + Optional entry = keyDictionaryRepository.findById(id); + + assertThat(entry.isPresent()).isTrue(); + assertThat(entry.get().getKeyId()).isEqualTo(first); + + keyDictionaryRepository.deleteById(id); + } + +} From d66e9ecf74b430ac6ab9bee4fa6b863157fbb419 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 8 Dec 2025 17:45:00 +0200 Subject: [PATCH 5/8] Added new lines to the end of files --- .../server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java | 2 +- .../server/dao/sqlts/dictionary/KeyDictionaryRepository.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java index d98105c0bc..8365e8facc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java @@ -39,4 +39,4 @@ public final class KeyDictionaryEntry { @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int", insertable = false, updatable = false) private Integer keyId; -} \ No newline at end of file +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java index a1141b62ff..e836cedb19 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java @@ -35,4 +35,4 @@ public interface KeyDictionaryRepository extends JpaRepository Date: Mon, 8 Dec 2025 18:36:01 +0200 Subject: [PATCH 6/8] Added find before lock for the startup of application --- .../dao/sqlts/dictionary/JpaKeyDictionaryDao.java | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index 1ef2c2d6aa..bcb9284371 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -50,6 +50,15 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService if (cached != null) { return cached; } + var compositeKey = new KeyDictionaryCompositeKey(strKey); + Optional entryOpt = keyDictionaryRepository.findById(compositeKey); + if (entryOpt.isPresent()) { + Integer keyId = entryOpt.get().getKeyId(); + if (keyId != null) { + keyDictionaryMap.put(strKey, keyId); + return keyId; + } + } creationLock.lock(); try { Integer keyId = keyDictionaryMap.get(strKey); @@ -59,8 +68,7 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); if (keyId == null || keyId == 0) { log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); - KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(strKey); - Optional entryOpt = keyDictionaryRepository.findById(id); + entryOpt = keyDictionaryRepository.findById(compositeKey); if (entryOpt.isEmpty() || entryOpt.get().getKeyId() == null || entryOpt.get().getKeyId() == 0) { From c0af057590273148da4e65fdb147ba3a9c38a23b Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 9 Dec 2025 11:31:02 +0200 Subject: [PATCH 7/8] refactoring logic in getOrSaveKeyId method --- .../sqlts/dictionary/JpaKeyDictionaryDao.java | 44 ++++++++++--------- 1 file changed, 23 insertions(+), 21 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index bcb9284371..9d23bdbcb5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -51,32 +51,29 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return cached; } var compositeKey = new KeyDictionaryCompositeKey(strKey); - Optional entryOpt = keyDictionaryRepository.findById(compositeKey); - if (entryOpt.isPresent()) { - Integer keyId = entryOpt.get().getKeyId(); - if (keyId != null) { - keyDictionaryMap.put(strKey, keyId); - return keyId; - } + Optional existingId = keyDictionaryRepository.findById(compositeKey) + .map(KeyDictionaryEntry::getKeyId) + .filter(id -> id != 0); + if (existingId.isPresent()) { + return cacheAndReturn(strKey, existingId.get()); } creationLock.lock(); try { - Integer keyId = keyDictionaryMap.get(strKey); - if (keyId != null) { - return keyId; + Integer fromCache = keyDictionaryMap.get(strKey); + if (fromCache != null) { + return fromCache; } - keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); - if (keyId == null || keyId == 0) { - log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); - entryOpt = keyDictionaryRepository.findById(compositeKey); - if (entryOpt.isEmpty() || - entryOpt.get().getKeyId() == null || - entryOpt.get().getKeyId() == 0) { - throw new IllegalStateException("Failed to resolve keyId for string key: " + strKey + " after fallback. keyId: " + keyId); - } + Integer keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); + if (keyId != null && keyId != 0) { + return cacheAndReturn(strKey, keyId); } - keyDictionaryMap.put(strKey, keyId); - return keyId; + log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); + keyId = keyDictionaryRepository.findById(compositeKey) + .map(KeyDictionaryEntry::getKeyId) + .filter(id -> id != 0) + .orElseThrow(() -> new IllegalStateException( + "Failed to resolve keyId for string key: " + strKey + " after fallback.")); + return cacheAndReturn(strKey, keyId); } finally { creationLock.unlock(); } @@ -93,4 +90,9 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return DaoUtil.pageToPageData(keyDictionaryRepository.findAll(DaoUtil.toPageable(pageLink))); } + private Integer cacheAndReturn(String key, Integer keyId) { + keyDictionaryMap.put(key, keyId); + return keyId; + } + } From 2c04943984abebdd7fe07772046e20d2173e8531 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 9 Dec 2025 12:11:14 +0200 Subject: [PATCH 8/8] removed checks for 0 value --- .../server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index 9d23bdbcb5..46bc10010e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -51,9 +51,7 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return cached; } var compositeKey = new KeyDictionaryCompositeKey(strKey); - Optional existingId = keyDictionaryRepository.findById(compositeKey) - .map(KeyDictionaryEntry::getKeyId) - .filter(id -> id != 0); + Optional existingId = keyDictionaryRepository.findById(compositeKey).map(KeyDictionaryEntry::getKeyId); if (existingId.isPresent()) { return cacheAndReturn(strKey, existingId.get()); } @@ -64,13 +62,12 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return fromCache; } Integer keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); - if (keyId != null && keyId != 0) { + if (keyId != null) { return cacheAndReturn(strKey, keyId); } log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); keyId = keyDictionaryRepository.findById(compositeKey) .map(KeyDictionaryEntry::getKeyId) - .filter(id -> id != 0) .orElseThrow(() -> new IllegalStateException( "Failed to resolve keyId for string key: " + strKey + " after fallback.")); return cacheAndReturn(strKey, keyId);