From 7fcd948071ffbefb0ff654758e58346a05150546 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 6 Feb 2025 16:11:58 +0200 Subject: [PATCH] fixed error when no telemetry in db --- .../cf/DefaultCalculatedFieldCache.java | 3 ++- ...efaultCalculatedFieldExecutionService.java | 5 ++++- .../ctx/state/BaseCalculatedFieldState.java | 2 +- .../cf/ctx/state/RocksDBStateService.java | 19 +++++++++++++++---- .../ctx/state/SingleValueArgumentEntry.java | 9 ++++++++- .../state/SingleValueArgumentEntryTest.java | 4 ++++ .../server/common/data/event/EventFilter.java | 3 ++- 7 files changed, 36 insertions(+), 9 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java index 9688f35fef..8fffa0029c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -32,8 +32,8 @@ import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.msg.cf.CalculatedFieldInitMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldLinkInitMsg; import org.thingsboard.server.dao.cf.CalculatedFieldService; -import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.dao.usagerecord.ApiLimitService; +import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import java.util.Collections; @@ -117,6 +117,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { CalculatedField calculatedField = getCalculatedField(calculatedFieldId); if (calculatedField != null) { ctx = new CalculatedFieldCtx(calculatedField, tbelInvokeService, apiLimitService); + ctx.init(); calculatedFieldsCtx.put(calculatedFieldId, ctx); log.debug("[{}] Put calculated field ctx into cache: {}", calculatedFieldId, ctx); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index 103997f0c7..43e85d87e9 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -39,6 +39,7 @@ import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.OutputType; @@ -68,7 +69,6 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.gen.transport.TransportProtos.AttributeScopeProto; @@ -447,6 +447,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas 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)); } 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 86d83a1f70..f1ce8038c7 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 @@ -52,7 +52,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { ArgumentEntry newEntry = entry.getValue(); ArgumentEntry existingEntry = arguments.get(key); - if (existingEntry == null) { + if (existingEntry == null || existingEntry == SingleValueArgumentEntry.EMPTY || existingEntry == TsRollingArgumentEntry.EMPTY) { validateNewEntry(newEntry); arguments.put(key, newEntry); stateUpdated = true; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java index 8a6a5c9cb7..b2e33e1705 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java @@ -63,7 +63,7 @@ public class RocksDBStateService implements CalculatedFieldStateService { CalculatedFieldStateProto stateProto = toProto(stateId, state); long maxStateSizeInKBytes = ctx.getMaxStateSizeInKBytes(); if (maxStateSizeInKBytes <= 0 || stateProto.getSerializedSize() <= ctx.getMaxStateSizeInKBytes()) { - rocksDBService.put(toProto(stateId), toProto(stateId, state)); + rocksDBService.put(toProto(stateId), stateProto); } callback.onSuccess(); } @@ -111,8 +111,11 @@ public class RocksDBStateService implements CalculatedFieldStateService { private SingleValueArgumentProto toSingleValueArgumentProto(String argName, SingleValueArgumentEntry entry) { SingleValueArgumentProto.Builder builder = SingleValueArgumentProto.newBuilder() - .setArgName(argName) - .setValue(KvProtoUtil.toTsValueProto(entry.getTs(), entry.getKvEntryValue())); + .setArgName(argName); + + if (entry != SingleValueArgumentEntry.EMPTY) { + builder.setValue(KvProtoUtil.toTsValueProto(entry.getTs(), entry.getKvEntryValue())); + } Optional.ofNullable(entry.getVersion()).ifPresent(builder::setVersion); @@ -122,7 +125,9 @@ public class RocksDBStateService implements CalculatedFieldStateService { private TsValueListProto toRollingArgumentProto(String argName, TsRollingArgumentEntry entry) { TsValueListProto.Builder builder = TsValueListProto.newBuilder().setKey(argName); - entry.getTsRecords().forEach((ts, value) -> builder.addTsValue(KvProtoUtil.toTsValueProto(ts, value))); + if (entry != TsRollingArgumentEntry.EMPTY) { + entry.getTsRecords().forEach((ts, value) -> builder.addTsValue(KvProtoUtil.toTsValueProto(ts, value))); + } return builder.build(); } @@ -151,6 +156,9 @@ public class RocksDBStateService implements CalculatedFieldStateService { } private SingleValueArgumentEntry fromSingleValueArgumentProto(SingleValueArgumentProto proto) { + if (!proto.hasValue()) { + return (SingleValueArgumentEntry) SingleValueArgumentEntry.EMPTY; + } TsValueProto tsValueProto = proto.getValue(); long ts = tsValueProto.getTs(); BasicKvEntry kvEntry = (BasicKvEntry) KvProtoUtil.fromTsValueProto(proto.getArgName(), tsValueProto); @@ -158,6 +166,9 @@ public class RocksDBStateService implements CalculatedFieldStateService { } private TsRollingArgumentEntry fromRollingArgumentProto(TsValueListProto proto) { + if (proto.getTsValueCount() <= 0) { + return (TsRollingArgumentEntry) TsRollingArgumentEntry.EMPTY; + } TreeMap tsRecords = new TreeMap<>(); proto.getTsValueList().forEach(tsValueProto -> { BasicKvEntry kvEntry = (BasicKvEntry) KvProtoUtil.fromTsValueProto(proto.getKey(), tsValueProto); 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 0832e53e5e..d7e5ddd017 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 @@ -94,10 +94,17 @@ public class SingleValueArgumentEntry implements ArgumentEntry { Long newVersion = singleValueEntry.getVersion(); if (newVersion == null || this.version == null || newVersion > this.version) { this.ts = singleValueEntry.getTs(); - this.kvEntryValue = singleValueEntry.getKvEntryValue(); this.version = newVersion; + + // TODO: should we persist updated ts and version values? + BasicKvEntry newValue = singleValueEntry.getKvEntryValue(); + if (this.kvEntryValue.getValue().equals(newValue.getValue())) { + return false; + } + this.kvEntryValue = singleValueEntry.getKvEntryValue(); return true; } + } else { throw new IllegalArgumentException("Unsupported argument entry type for single value argument entry: " + entry.getType()); } 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 203d7b3d71..13651e852d 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 @@ -69,4 +69,8 @@ public class SingleValueArgumentEntryTest { assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 18L), 234L))).isFalse(); } + @Test + void testUpdateEntryWhenValueWasNotChanged() { + assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 11L), 237L))).isFalse(); + } } \ No newline at end of file diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventFilter.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventFilter.java index 6d2a110cf7..4c9791e3fb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventFilter.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventFilter.java @@ -29,7 +29,8 @@ import io.swagger.v3.oas.annotations.media.Schema; @JsonSubTypes.Type(value = RuleChainDebugEventFilter.class, name = "DEBUG_RULE_CHAIN"), @JsonSubTypes.Type(value = ErrorEventFilter.class, name = "ERROR"), @JsonSubTypes.Type(value = LifeCycleEventFilter.class, name = "LC_EVENT"), - @JsonSubTypes.Type(value = StatisticsEventFilter.class, name = "STATS") + @JsonSubTypes.Type(value = StatisticsEventFilter.class, name = "STATS"), + @JsonSubTypes.Type(value = CalculatedFieldDebugEventFilter.class, name = "DEBUG_CALCULATED_FIELD") }) public interface EventFilter {