Browse Source

Merge pull request #14531 from irynamatveieva/master-fix/cf-args

Fixed timestamp handling for calculated field arguments with missing telemetry
pull/14542/head
Viacheslav Klimov 10 months ago
committed by GitHub
parent
commit
ca7703d44a
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 43
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java
  2. 22
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/HasLatestTs.java
  3. 9
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java
  4. 18
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java
  5. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java
  6. 1
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java
  7. 85
      application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java

43
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java

@ -26,8 +26,6 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.aggregation.single.EntityAggregationArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingZoneState;
import org.thingsboard.server.utils.CalculatedFieldUtils;
import java.io.Closeable;
@ -41,7 +39,7 @@ import java.util.stream.Collectors;
@Getter
public abstract class BaseCalculatedFieldState implements CalculatedFieldState, Closeable {
protected static final long DEFAULT_LAST_UPDATE_TS = -1L;
public static final long DEFAULT_LAST_UPDATE_TS = -1L;
protected final EntityId entityId;
protected CalculatedFieldCtx ctx;
@ -103,7 +101,6 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
updatedArguments = new HashMap<>(argumentValues.size());
}
updatedArguments.put(key, newEntry);
updateLastUpdateTimestamp(newEntry);
}
}
@ -161,23 +158,29 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
return resultNode;
}
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<Long, Double> lastEntry = tsRollingArgumentEntry.getTsRecords().lastEntry();
newTs = (lastEntry != null) ? lastEntry.getKey() : DEFAULT_LAST_UPDATE_TS;
} else if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) {
newTs = relatedEntitiesArgumentEntry.getEntityInputs().values().stream()
.mapToLong(e -> (e instanceof SingleValueArgumentEntry s) ? s.getTs() : DEFAULT_LAST_UPDATE_TS)
.max()
.orElse(DEFAULT_LAST_UPDATE_TS);
} else if (entry instanceof GeofencingArgumentEntry geofencingArgumentEntry) {
newTs = geofencingArgumentEntry.getZoneStates().values().stream()
.mapToLong(GeofencingZoneState::getTs).max().orElse(DEFAULT_LAST_UPDATE_TS);
public long getLatestTimestamp() {
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 HasLatestTs hasLatestTsEntry) {
latestTs = Math.max(latestTs, hasLatestTsEntry.getLatestTs());
}
}
this.latestTimestamp = Math.max(this.latestTimestamp, newTs);
return latestTs;
}
protected ReadinessStatus checkReadiness(List<String> requiredArguments, Map<String, ArgumentEntry> currentArguments) {

22
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/HasLatestTs.java

@ -0,0 +1,22 @@
/**
* 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.service.cf.ctx.state;
public interface HasLatestTs {
long getLatestTs();
}

9
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java

@ -31,11 +31,13 @@ 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
@Slf4j
public class TsRollingArgumentEntry implements ArgumentEntry {
public class TsRollingArgumentEntry implements ArgumentEntry, HasLatestTs {
private Integer limit;
private Long timeWindow;
@ -83,6 +85,11 @@ public class TsRollingArgumentEntry implements ArgumentEntry {
return tsRecords;
}
public long getLatestTs() {
var lastEntry = tsRecords.lastEntry();
return (lastEntry != null) ? lastEntry.getKey() : DEFAULT_LAST_UPDATE_TS;
}
@Override
public TbelCfArg toTbelCfArg() {
List<TbelCfTsDoubleVal> values = new ArrayList<>(tsRecords.size());

18
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesArgumentEntry.java

@ -23,14 +23,17 @@ import org.thingsboard.script.api.tbel.TbelCfSingleValueArg;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType;
import org.thingsboard.server.service.cf.ctx.state.HasLatestTs;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import java.util.Map;
import java.util.stream.Collectors;
import static org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState.DEFAULT_LAST_UPDATE_TS;
@Data
@AllArgsConstructor
public class RelatedEntitiesArgumentEntry implements ArgumentEntry {
public class RelatedEntitiesArgumentEntry implements ArgumentEntry, HasLatestTs {
private final Map<EntityId, ArgumentEntry> entityInputs;
@ -46,6 +49,19 @@ public class RelatedEntitiesArgumentEntry implements ArgumentEntry {
return entityInputs;
}
@Override
public long getLatestTs() {
long latestTs = DEFAULT_LAST_UPDATE_TS;
for (ArgumentEntry entry : entityInputs.values()) {
if (entry instanceof SingleValueArgumentEntry single) {
if (!single.isDefaultValue()) {
latestTs = Math.max(latestTs, single.getTs());
}
}
}
return latestTs;
}
@Override
public boolean updateEntry(ArgumentEntry entry) {
if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) {

11
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java

@ -25,13 +25,16 @@ import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType;
import org.thingsboard.server.service.cf.ctx.state.HasLatestTs;
import java.util.Map;
import java.util.stream.Collectors;
import static org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState.DEFAULT_LAST_UPDATE_TS;
@Data
@Slf4j
public class GeofencingArgumentEntry implements ArgumentEntry {
public class GeofencingArgumentEntry implements ArgumentEntry, HasLatestTs {
private Map<EntityId, GeofencingZoneState> zoneStates;
@ -58,6 +61,12 @@ public class GeofencingArgumentEntry implements ArgumentEntry {
return zoneStates;
}
@Override
public long getLatestTs() {
return zoneStates.values().stream()
.mapToLong(GeofencingZoneState::getTs).max().orElse(DEFAULT_LAST_UPDATE_TS);
}
@Override
public boolean updateEntry(ArgumentEntry entry) {
if (!(entry instanceof GeofencingArgumentEntry geofencingArgumentEntry)) {

1
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java

@ -104,7 +104,6 @@ public class CalculatedFieldArgumentUtils {
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);
}

85
application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java

@ -580,6 +580,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes
@Test
public void testScriptCalculatedFieldWhenUsedLatestTsInScript() throws Exception {
Device testDevice = createDevice("Test device", "1234567890");
long ts = System.currentTimeMillis() - 300000L;
postTelemetry(testDevice.getId(), String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts));
@ -614,6 +615,90 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes
});
}
@Test
public void testSimpleCalculatedFieldWhenUseLatestTsIsTrueAndDefaultArguments() throws Exception {
Device testDevice = createDevice("Test device", "1234567890");
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");
TimeSeriesOutput output = new TimeSeriesOutput();
output.setName("d");
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 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");
});
}
@Test
public void testSimpleCalculatedFieldWhenCtxBecameUninitialized() throws Exception {
Device testDevice = createDevice("Test device", "1234567890");

Loading…
Cancel
Save