Browse Source

Merge branch 'rc' into features/add_tooltip_option_to_show_stack_mode_total_value_on_timeseries_chart_widgets

pull/13373/head
Paolo Cristiani 1 year ago
committed by GitHub
parent
commit
2acc2bad5c
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 39
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntryTest.java
  2. 1
      common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java
  3. 1
      common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java
  4. 1
      common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java
  5. 51
      common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java
  6. 1
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java

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

@ -17,8 +17,15 @@ package org.thingsboard.server.service.cf.ctx.state;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.thingsboard.script.api.tbel.TbelCfArg;
import org.thingsboard.script.api.tbel.TbelCfSingleValueArg;
import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@ -73,4 +80,34 @@ public class SingleValueArgumentEntryTest {
void testUpdateEntryWhenValueWasNotChanged() {
assertThat(entry.updateEntry(new SingleValueArgumentEntry(ts + 18, new LongDataEntry("key", 11L), 364L))).isTrue();
}
}
@Test
void testToTbelCfArgWhenJsonIsObject() {
entry = new SingleValueArgumentEntry(ts, new JsonDataEntry("key", "{\"test\": 10}"), 370L);
TbelCfArg tbelCfArg = entry.toTbelCfArg();
assertThat(tbelCfArg).isNotNull();
assertThat(tbelCfArg).isInstanceOf(TbelCfSingleValueArg.class);
TbelCfSingleValueArg singleValueArg = (TbelCfSingleValueArg) tbelCfArg;
assertThat(singleValueArg.getValue()).isInstanceOf(Map.class);
Map<String, Integer> expectedMap = Map.of("test", 10);
assertThat(singleValueArg.getValue()).isEqualTo(expectedMap);
}
@Test
void testToTbelCfArgWhenJsonIsArray() {
entry = new SingleValueArgumentEntry(ts, new JsonDataEntry("key", "[{\"test\": 10}, {\"test2\": 20}]"), 371L);
TbelCfArg tbelCfArg = entry.toTbelCfArg();
assertThat(tbelCfArg).isNotNull();
assertThat(tbelCfArg).isInstanceOf(TbelCfSingleValueArg.class);
TbelCfSingleValueArg singleValueArg = (TbelCfSingleValueArg) tbelCfArg;
assertThat(singleValueArg.getValue()).isInstanceOf(List.class);
List<Map<String, Integer>> expectedList = new ArrayList<>();
expectedList.add(Map.of("test", 10));
expectedList.add(Map.of("test2", 20));
assertThat(singleValueArg.getValue()).isEqualTo(expectedList);
}
}

1
common/data/src/main/java/org/thingsboard/server/common/data/limit/LimitedApi.java

@ -43,6 +43,7 @@ public enum LimitedApi {
RateLimitUtil.merge(
DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantCoreRateLimits,
DefaultTenantProfileConfiguration::getCassandraWriteQueryTenantRuleEngineRateLimits), "Monolith telemetry Cassandra write queries", true),
CASSANDRA_QUERIES(null, true), // left for backward compatibility with RateLimitsNotificationInfo
EDGE_EVENTS(DefaultTenantProfileConfiguration::getEdgeEventRateLimits, "Edge events", true),
EDGE_EVENTS_PER_EDGE(DefaultTenantProfileConfiguration::getEdgeEventRateLimitsPerEdge, "Edge events per edge", false),
EDGE_UPLINK_MESSAGES(DefaultTenantProfileConfiguration::getEdgeUplinkMessagesRateLimits, "Edge uplink messages", true),

1
common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java

@ -277,6 +277,7 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
eventConsumer.awaitStop();
responseTemplate.stop();
stateService.stop();
versionsStore.shutdown();
}
}

1
common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java

@ -224,6 +224,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
stateConsumer.awaitStop();
eventsToBackupConsumer.stop();
stateProducer.stop();
versionsStore.shutdown();
}
}

51
common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java

@ -15,31 +15,35 @@
*/
package org.thingsboard.server.edqs.util;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.edqs.EdqsObjectKey;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@Slf4j
public class VersionsStore {
private final Cache<EdqsObjectKey, Long> versions;
private final ConcurrentMap<EdqsObjectKey, TimedValue> versions = new ConcurrentHashMap<>();
private final long expirationMillis;
private final ScheduledExecutorService cleaner = Executors.newSingleThreadScheduledExecutor();
public VersionsStore(int ttlMinutes) {
this.versions = Caffeine.newBuilder()
.expireAfterWrite(ttlMinutes, TimeUnit.MINUTES)
.build();
this.expirationMillis = TimeUnit.MINUTES.toMillis(ttlMinutes);
startCleanupTask();
}
public boolean isNew(EdqsObjectKey key, Long version) {
AtomicBoolean isNew = new AtomicBoolean(false);
versions.asMap().compute(key, (k, prevVersion) -> {
if (prevVersion == null || prevVersion <= version) {
versions.compute(key, (k, prevVersion) -> {
if (prevVersion == null || prevVersion.value <= version) {
isNew.set(true);
return version;
return new TimedValue(version);
} else {
log.debug("[{}] Version {} is outdated, the latest is {}", key, version, prevVersion);
return prevVersion;
@ -48,4 +52,33 @@ public class VersionsStore {
return isNew.get();
}
private void startCleanupTask() {
cleaner.scheduleAtFixedRate(() -> {
try {
long now = System.currentTimeMillis();
for (Map.Entry<EdqsObjectKey, TimedValue> entry : versions.entrySet()) {
if (now - entry.getValue().lastUpdated > expirationMillis) {
versions.remove(entry.getKey(), entry.getValue());
}
}
} catch (Exception e) {
log.error("Cleanup task failed", e);
}
}, expirationMillis, expirationMillis, TimeUnit.MILLISECONDS);
}
public void shutdown() {
cleaner.shutdown();
}
private static class TimedValue {
private final long lastUpdated;
private final long value;
public TimedValue(long value) {
this.value = value;
this.lastUpdated = System.currentTimeMillis();
}
}
}

1
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java

@ -33,7 +33,6 @@ import org.thingsboard.server.common.data.notification.rule.DefaultNotificationR
import org.thingsboard.server.common.data.notification.rule.EscalatedNotificationRuleRecipientsConfig;
import org.thingsboard.server.common.data.notification.rule.NotificationRule;
import org.thingsboard.server.common.data.notification.rule.NotificationRuleConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.ResourcesShortageTrigger.Resource;
import org.thingsboard.server.common.data.notification.rule.trigger.config.AlarmAssignmentNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.config.AlarmCommentNotificationRuleTriggerConfig;
import org.thingsboard.server.common.data.notification.rule.trigger.config.AlarmNotificationRuleTriggerConfig;

Loading…
Cancel
Save