Browse Source

CF: fix update of entry with default value

pull/14193/head
VIacheslavKlimov 12 months ago
parent
commit
a12c0d0704
  1. 7
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  2. 13
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java
  4. 6
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java

7
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -43,6 +43,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; 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.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import java.util.HashMap; import java.util.HashMap;
import java.util.List; import java.util.List;
@ -226,12 +227,12 @@ public abstract class AbstractCalculatedFieldProcessingService {
return Futures.transform(attributeOptFuture, attrOpt -> { return Futures.transform(attributeOptFuture, attrOpt -> {
log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt); log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt);
AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, 0L)); AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, SingleValueArgumentEntry.DEFAULT_VERSION));
return transformSingleValueArgument(Optional.of(attributeKvEntry)); return transformSingleValueArgument(Optional.of(attributeKvEntry));
}, calculatedFieldCallbackExecutor); }, calculatedFieldCallbackExecutor);
} }
protected ListenableFuture<ArgumentEntry> fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { protected ListenableFuture<ArgumentEntry> fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long defaultTs) {
String timeseriesKey = argument.getRefEntityKey().getKey(); String timeseriesKey = argument.getRefEntityKey().getKey();
log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, timeseriesKey); log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, timeseriesKey);
return transformSingleValueArgument( return transformSingleValueArgument(
@ -239,7 +240,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
timeseriesService.findLatest(tenantId, entityId, timeseriesKey), timeseriesService.findLatest(tenantId, entityId, timeseriesKey),
result -> { result -> {
log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, timeseriesKey, result); log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, timeseriesKey, result);
return result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))); return result.or(() -> Optional.of(new BasicTsKvEntry(defaultTs, createDefaultKvEntry(argument), SingleValueArgumentEntry.DEFAULT_VERSION)));
}, calculatedFieldCallbackExecutor)); }, calculatedFieldCallbackExecutor));
} }

13
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java

@ -43,6 +43,8 @@ public class SingleValueArgumentEntry implements ArgumentEntry {
private boolean forceResetPrevious; private boolean forceResetPrevious;
public static final Long DEFAULT_VERSION = -1L;
public SingleValueArgumentEntry(TsKvProto entry) { public SingleValueArgumentEntry(TsKvProto entry) {
this.ts = entry.getTs(); this.ts = entry.getTs();
if (entry.hasVersion()) { if (entry.hasVersion()) {
@ -112,8 +114,10 @@ public class SingleValueArgumentEntry implements ArgumentEntry {
@Override @Override
public boolean updateEntry(ArgumentEntry entry) { public boolean updateEntry(ArgumentEntry entry) {
if (entry instanceof SingleValueArgumentEntry singleValueEntry) { if (entry instanceof SingleValueArgumentEntry singleValueEntry) {
if (singleValueEntry.getTs() <= this.ts) { if (singleValueEntry.getTs() < this.ts) {
return false; if (!isDefaultValue()) {
return false;
}
} }
Long newVersion = singleValueEntry.getVersion(); Long newVersion = singleValueEntry.getVersion();
@ -128,4 +132,9 @@ public class SingleValueArgumentEntry implements ArgumentEntry {
} }
return false; return false;
} }
public boolean isDefaultValue() {
return DEFAULT_VERSION.equals(this.version);
}
} }

2
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java

@ -300,8 +300,6 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
private TbAlarmResult calculateAlarmResult(AlarmRuleState ruleState, CalculatedFieldCtx ctx) { private TbAlarmResult calculateAlarmResult(AlarmRuleState ruleState, CalculatedFieldCtx ctx) {
AlarmSeverity severity = ruleState.getSeverity(); AlarmSeverity severity = ruleState.getSeverity();
if (currentAlarm != null) { if (currentAlarm != null) {
// TODO: In some extremely rare cases, we might miss the event of alarm clear (If one use in-mem queue and restarted the server) or (if one manipulated the rule chain).
// Maybe we should fetch alarm every time?
currentAlarm.setEndTs(System.currentTimeMillis()); currentAlarm.setEndTs(System.currentTimeMillis());
AlarmSeverity oldSeverity = currentAlarm.getSeverity(); AlarmSeverity oldSeverity = currentAlarm.getSeverity();
// Skip update if severity is decreased. // Skip update if severity is decreased.

6
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java

@ -57,6 +57,11 @@ public class SingleValueArgumentEntryTest {
assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 363L))).isFalse(); assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 363L))).isFalse();
} }
@Test
void testUpdateEntryWithTheSameTsAndDifferentVersion() {
assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts, new LongDataEntry("key", 13L), 364L))).isTrue();
}
@Test @Test
void testUpdateEntryWhenNewVersionIsNull() { void testUpdateEntryWhenNewVersionIsNull() {
assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 16, new LongDataEntry("key", 13L), null))).isTrue(); assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 16, new LongDataEntry("key", 13L), null))).isTrue();
@ -115,4 +120,5 @@ public class SingleValueArgumentEntryTest {
expectedList.add(Map.of("test2", 20)); expectedList.add(Map.of("test2", 20));
assertThat(singleValueArg.getValue()).isEqualTo(expectedList); assertThat(singleValueArg.getValue()).isEqualTo(expectedList);
} }
} }

Loading…
Cancel
Save