Browse Source

Simplify tests

pull/14549/head
Volodymyr Babak 9 months ago
parent
commit
c42705c032
  1. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java
  3. 3
      application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java
  4. 237
      application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java
  5. 88
      application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java
  6. 3
      dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java

6
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -546,7 +546,8 @@ public abstract class EdgeGrpcSession implements Closeable {
try {
if (msg.getSuccess()) {
sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId());
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PUSHED, edge.getTenantId(), edge.getId(), 1));
ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PUSHED, edge.getTenantId(), edge.getId(), 1));
log.debug("[{}][{}][{}] Msg has been processed successfully! Msg Id: [{}], Msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg);
} else {
log.debug("[{}][{}][{}] Msg processing failed! Msg Id: [{}], Error msg: {}", tenantId, edge.getId(), sessionId, msg.getDownlinkMsgId(), msg.getErrorMsg());
@ -815,7 +816,8 @@ public abstract class EdgeGrpcSession implements Closeable {
}
}
highPriorityQueue.add(edgeEvent);
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edge.getTenantId(), edgeEvent.getEdgeId(), 1));
ctx.getStatsCounterService().ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edge.getTenantId(), edgeEvent.getEdgeId(), 1));
}
protected ListenableFuture<List<Void>> processUplinkMsg(UplinkMsg uplinkMsg) {

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java

@ -52,7 +52,8 @@ public class KafkaEdgeEventService extends BaseEdgeEventService {
TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId());
ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build();
producerProvider.getTbEdgeEventsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null);
statsCounterService.ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
return Futures.immediateFuture(null);
}

3
application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java

@ -54,7 +54,7 @@ import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_P
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_TMP_FAILED;
@TbCoreComponent
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true", matchIfMissing = false)
@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true")
@RequiredArgsConstructor
@Service
@Slf4j
@ -70,7 +70,6 @@ public class EdgeStatsService {
@Value("${edges.stats.report-interval-millis:600000}")
private long reportIntervalMillis;
@Scheduled(
fixedDelayString = "${edges.stats.report-interval-millis:600000}",
initialDelayString = "${edges.stats.report-interval-millis:600000}"

237
application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java

@ -16,8 +16,6 @@
package org.thingsboard.server.edge;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.protobuf.AbstractMessage;
import lombok.extern.slf4j.Slf4j;
import org.junit.Assert;
import org.junit.Test;
@ -35,24 +33,16 @@ import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService;
import org.thingsboard.server.dao.edge.stats.EdgeStatsKey;
import org.thingsboard.server.dao.edge.stats.MsgCounters;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.edge.v1.EntityDataProto;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.edge.stats.EdgeStatsService;
import java.time.Duration;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import static org.awaitility.Awaitility.await;
import static org.junit.jupiter.api.Assertions.assertAll;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_ADDED;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED;
import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_PUSHED;
@ -64,19 +54,10 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
private static final String STATISTICS_DEVICE_PROFILE = "STATISTICS";
private final Map<EdgeStatsKey, Long> EXPECTED_EDGE_STATS = Map.of(
DOWNLINK_MSGS_ADDED, 6L,
DOWNLINK_MSGS_PUSHED, 6L,
DOWNLINK_MSGS_PERMANENTLY_FAILED, 0L,
DOWNLINK_MSGS_TMP_FAILED, 0L
);
private final Map<EdgeStatsKey, Long> EXPECTED_EMPTY_EDGE_STATS = Map.of(
DOWNLINK_MSGS_ADDED, 0L,
DOWNLINK_MSGS_PUSHED, 0L,
DOWNLINK_MSGS_PERMANENTLY_FAILED, 0L,
DOWNLINK_MSGS_TMP_FAILED, 0L
);
private static final long EXPECTED_MSGS_ADDED = 6L;
private static final long EXPECTED_MSGS_PUSHED = 6L;
private static final long EXPECTED_MSGS_PERMANENTLY_FAILED = 0L;
private static final long EXPECTED_MSGS_TMP_FAILED = 0L;
@Autowired
private EdgeStatsService edgeStatsService;
@ -85,209 +66,83 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest {
@Test
public void testFullEdgeStatsCycle() throws Exception {
// 1. Clear previous stats and prepare test data
// GIVEN
prepareTestData();
// 2. Wait until Edge counters are updated
awaitEdgeCountersUpdated();
// 3. Report statistics
edgeStatsService.reportStats();
// 4. Wait until timeseries data is persisted
awaitStatsMatch(EXPECTED_EDGE_STATS);
// 5. Send statistics to Edge
List<TsKvEntry> latestStatsEntries = tsService.findLatest(
tenantId,
edge.getId(),
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()
).get();
sendStatsToEdge(latestStatsEntries);
// 6. Wait until telemetry Proto contains the expected stats
// Await Edge Counters Updated
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
EntityDataProto latestMsg = getLatestEntityDataMessage();
Map<String, Long> actualStats = toMap(latestMsg.getPostTelemetryMsg().getTsKvList(0));
assertAllStatsEqual(EXPECTED_EDGE_STATS, actualStats, "Proto stats");
MsgCounters counters = statsCounterService.getMsgCountersByEdge().get(edge.getId());
assertEquals(EXPECTED_MSGS_ADDED, counters.getMsgsAdded().get());
assertEquals(EXPECTED_MSGS_PUSHED, counters.getMsgsPushed().get());
assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, counters.getMsgsPermanentlyFailed().get());
assertEquals(EXPECTED_MSGS_TMP_FAILED, counters.getMsgsTmpFailed().get());
});
}
@Test
public void testNoMessagesFromEdge() throws ExecutionException, InterruptedException {
// 1. Clear stats counters for the Edge
statsCounterService.clear(edge.getId());
Thread.sleep(1000);
// 2. Report stats with no data from the Edge
// WHEN
edgeStatsService.reportStats();
// 3. Verify that persisted timeseries contains only empty stats
awaitStatsMatch(EXPECTED_EMPTY_EDGE_STATS);
List<TsKvEntry> actual = tsService.findLatest(
tenantId,
edge.getId(),
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()
).get();
assertAllStatsEqual(EXPECTED_EMPTY_EDGE_STATS, toMap(actual), "Empty stats");
}
@Test
public void testRepeatedReportStatsDoesNotDuplicate() throws ExecutionException, InterruptedException {
// 1. Clear previous stats and prepare test data
prepareTestData();
// 2. Wait until Edge counters are updated
awaitEdgeCountersUpdated();
// 3. First report call
edgeStatsService.reportStats();
// 4. Verify that the persisted stats match expectations
awaitStatsMatch(EXPECTED_EDGE_STATS);
// 5. Remove persisted stats to simulate a re-report scenario
tsService.removeLatest(
tenantId,
edge.getId(),
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()
);
// 6. Second report call without increments (counters already cleared)
edgeStatsService.reportStats();
// 7. Verify that the stats are empty after the second report
awaitStatsMatch(EXPECTED_EMPTY_EDGE_STATS);
}
private void awaitStatsMatch(Map<EdgeStatsKey, Long> expected) {
// THEN
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
Map<String, Long> actualStats = fetchLatestStats();
assertAllStatsEqual(expected, actualStats, "Timeseries stats");
List<TsKvEntry> actualStats = fetchLatestStats();
assertEquals(EXPECTED_MSGS_ADDED, getStatsLongValue(actualStats, DOWNLINK_MSGS_ADDED));
assertEquals(EXPECTED_MSGS_PUSHED, getStatsLongValue(actualStats, DOWNLINK_MSGS_PUSHED));
assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, getStatsLongValue(actualStats, DOWNLINK_MSGS_PERMANENTLY_FAILED));
assertEquals(EXPECTED_MSGS_TMP_FAILED, getStatsLongValue(actualStats, DOWNLINK_MSGS_TMP_FAILED));
});
}
private Map<String, Long> fetchLatestStats() throws ExecutionException, InterruptedException {
List<TsKvEntry> latestStatsEntries = tsService.findLatest(
private long getStatsLongValue(List<TsKvEntry> stats, EdgeStatsKey key) {
return stats.stream().filter(e -> e.getKey().equals(key.getKey())).findFirst().get().getLongValue().orElse(0L);
}
private List<TsKvEntry> fetchLatestStats() throws ExecutionException, InterruptedException {
return tsService.findLatest(
tenantId,
edge.getId(),
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()
).get();
return toMap(latestStatsEntries);
Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()).get();
}
private void prepareTestData() throws InterruptedException, ExecutionException {
statsCounterService.clear(edge.getId());
// 2 stats messages: ADDED Device Profile, ASSIGN Device
// Save device and assign to edge
// 2 DOWNLINK_MSGS_ADDED, EdgeEvents: [{DEVICE_PROFILE: ADDED}, {DEVICE: ASSIGNED_TO_EDGE}]
// 2 DOWNLINK_MSGS_PUSHED, Downlinks: [{deviceProfileUpdateMsg}, {deviceUpdateMsg, deviceProfileUpdateMsg, deviceCredentialsUpdateMsg}]
edgeImitator.expectMessageAmount(4);
Device device = saveDevice("StatisticDevice", STATISTICS_DEVICE_PROFILE);
doPost("/api/edge/" + edge.getUuidId() + "/device/" + device.getUuidId(), Device.class);
edgeImitator.waitForMessages();
// 1 stats message: ASSIGN Asset
Device savedDevice = saveDevice("StatisticDevice", STATISTICS_DEVICE_PROFILE);
doPost("/api/edge/" + edge.getUuidId() + "/device/" + savedDevice.getUuidId(), Device.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// Save asset and assign to edge
// 1 DOWNLINK_MSGS_ADDED, EdgeEvents: [{ASSET: ASSIGNED_TO_EDGE}]
// 1 DOWNLINK_MSGS_PUSHED, Downlinks: [{assetUpdateMsg, assetProfileUpdateMsg}]
edgeImitator.expectMessageAmount(2);
Asset savedAsset = saveAsset("Edge Asset 2");
Asset savedAsset = saveAsset("Edge Asset");
doPost("/api/edge/" + edge.getUuidId()
+ "/asset/" + savedAsset.getUuidId(), Asset.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// 2 stats messages: ADDED Customer, ASSIGN Customer
edgeImitator.expectMessageAmount(1);
// Create customer and assign edge to the customer
// 2 DOWNLINK_MSGS_ADDED, EdgeEvents: [{CUSTOMER: ADDED}, {EDGE: ASSIGNED_TO_CUSTOMER}]
// 2 DOWNLINK_MSGS_PUSHED, Downlinks: [{customerUpdateMsg}, {edgeConfiguration}]
edgeImitator.expectMessageAmount(2);
Customer customer = new Customer();
customer.setTitle("Edge Customer");
Customer savedCustomer = doPost("/api/customer", customer, Customer.class);
Assert.assertFalse(edgeImitator.waitForMessages(5));
// assign edge to customer
edgeImitator.expectMessageAmount(2);
doPost("/api/customer/" + savedCustomer.getUuidId()
+ "/edge/" + edge.getUuidId(), Edge.class);
Assert.assertTrue(edgeImitator.waitForMessages());
// 1 stats message: Timeseries
// Send device telemetry downlink for the device
// 1 DOWNLINK_MSGS_ADDED, EdgeEvents: [{DEVICE: TIMESERIES_UPDATED}]
// 1 DOWNLINK_MSGS_PUSHED, Downlinks: [{entityData}]
edgeImitator.expectMessageAmount(1);
String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = JacksonUtil.toJsonNode(timeseriesData);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, savedDevice.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeEventService.saveAsync(edgeEvent).get();
Assert.assertTrue(edgeImitator.waitForMessages());
}
private void assertAllStatsEqual(Map<EdgeStatsKey, Long> expected, Map<String, Long> actual, String context) {
assertAll(context,
expected.entrySet().stream()
.map(e -> () -> assertEquals(e.getValue(), actual.get(e.getKey().getKey()), "Mismatch for stat: " + e.getKey()))
);
}
private void awaitEdgeCountersUpdated() {
await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> {
MsgCounters counters = statsCounterService.getMsgCountersByEdge().get(edge.getId());
Map<String, Long> actualCounters = toMap(
Map.entry(DOWNLINK_MSGS_ADDED.getKey(), () -> counters.getMsgsAdded().get()),
Map.entry(DOWNLINK_MSGS_PUSHED.getKey(), () -> counters.getMsgsPushed().get()),
Map.entry(DOWNLINK_MSGS_PERMANENTLY_FAILED.getKey(), () -> counters.getMsgsPermanentlyFailed().get()),
Map.entry(DOWNLINK_MSGS_TMP_FAILED.getKey(), () -> counters.getMsgsTmpFailed().get())
);
assertAllStatsEqual(EXPECTED_EDGE_STATS, actualCounters, "Edge counters");
});
}
@SafeVarargs
private Map<String, Long> toMap(Map.Entry<String, Supplier<Long>>... suppliers) {
return Arrays.stream(suppliers)
.collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get()));
}
private Map<String, Long> toMap(List<TsKvEntry> stats) {
return stats.stream()
.collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(0L)));
}
private Map<String, Long> toMap(TransportProtos.TsKvListProto kvList) {
Map<String, Long> map = kvList.getKvList().stream()
.collect(Collectors.toMap(
TransportProtos.KeyValueProto::getKey,
TransportProtos.KeyValueProto::getLongV
));
for (EdgeStatsKey key : EdgeStatsKey.values()) {
map.putIfAbsent(key.getKey(), 0L);
}
return map;
}
private void sendStatsToEdge(List<TsKvEntry> stats) throws Exception {
edgeImitator.expectMessageAmount(1);
EdgeEvent edgeEvent = constructEdgeEvent(
tenantId,
edge.getId(),
EdgeEventActionType.TIMESERIES_UPDATED,
edge.getId().getId(),
EdgeEventType.EDGE,
buildStatsJson(System.currentTimeMillis(), stats)
);
edgeEventService.saveAsync(edgeEvent).get();
assertTrue(edgeImitator.waitForMessages());
}
private EntityDataProto getLatestEntityDataMessage() {
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
assertInstanceOf(EntityDataProto.class, latestMessage);
EntityDataProto msg = (EntityDataProto) latestMessage;
assertEquals(edge.getUuidId().getMostSignificantBits(), msg.getEntityIdMSB());
assertEquals(edge.getUuidId().getLeastSignificantBits(), msg.getEntityIdLSB());
assertEquals(edge.getId().getEntityType().name(), msg.getEntityType());
assertTrue(msg.hasPostTelemetryMsg());
return msg;
}
private ObjectNode buildStatsJson(long ts, List<TsKvEntry> statsEntries) {
ObjectNode entityBody = JacksonUtil.newObjectNode();
entityBody.put("ts", ts);
ObjectNode data = JacksonUtil.newObjectNode();
statsEntries.forEach(entry -> data.put(entry.getKey(), entry.getValueAsString()));
entityBody.set("data", data);
return entityBody;
}
}

88
application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java

@ -64,6 +64,13 @@ public class EdgeStatsTest {
private static final int TTL_DAYS = 30;
private static final long REPORT_INTERVAL_MILLIS = 600_000L;
private static final long EXPECTED_MSGS_ADDED = 5L;
private static final long EXPECTED_MSGS_PUSHED = 3L;
private static final long EXPECTED_MSGS_PERMANENTLY_FAILED = 1L;
private static final long EXPECTED_MSGS_TMP_FAILED = 0L;
private static final long EXPECTED_MSGS_LAG = 10L;
private static final long EXPECTED_MSGS_KAFKA_LAG = 15L;
@Mock
private TimeseriesService tsService;
@Mock
@ -97,12 +104,38 @@ public class EdgeStatsTest {
@Test
public void testReportStatsSavesTelemetry() {
// GIVEN
setupCounters();
// WHEN
edgeStatsService.reportStats();
// THEN
Map<String, Long> counters = verifyCounters();
Assertions.assertEquals(EXPECTED_MSGS_LAG, counters.get(DOWNLINK_MSGS_LAG.getKey()).longValue());
}
@Test
public void testReportStatsWithKafkaLag() {
// GIVEN
setupCounters();
setupKafkaLag();
// WHEN
edgeStatsService.reportStats();
// THEN
Map<String, Long> valuesByKey = verifyCounters();
Assertions.assertEquals(EXPECTED_MSGS_KAFKA_LAG, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey()));
}
private void setupCounters() {
MsgCounters counters = new MsgCounters(tenantId);
counters.getMsgsAdded().set(5);
counters.getMsgsPushed().set(3);
counters.getMsgsPermanentlyFailed().set(1);
counters.getMsgsTmpFailed().set(0);
counters.getMsgsLag().set(10);
counters.getMsgsAdded().set(EXPECTED_MSGS_ADDED);
counters.getMsgsPushed().set(EXPECTED_MSGS_PUSHED);
counters.getMsgsPermanentlyFailed().set(EXPECTED_MSGS_PERMANENTLY_FAILED);
counters.getMsgsTmpFailed().set(EXPECTED_MSGS_TMP_FAILED);
counters.getMsgsLag().set(EXPECTED_MSGS_LAG);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters);
@ -111,9 +144,9 @@ public class EdgeStatsTest {
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));
}
edgeStatsService.reportStats();
private Map<String, Long> verifyCounters() {
verify(tsService, times(1)).save(eq(tenantId), eq(edgeId), anyList(), anyLong());
verify(statsCounterService, times(1)).clear(edgeId);
@ -123,50 +156,23 @@ public class EdgeStatsTest {
Map<String, Long> valuesByKey = entries.stream()
.collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L)));
Assertions.assertEquals(5L, valuesByKey.get(DOWNLINK_MSGS_ADDED.getKey()).longValue());
Assertions.assertEquals(3L, valuesByKey.get(DOWNLINK_MSGS_PUSHED.getKey()).longValue());
Assertions.assertEquals(1L, valuesByKey.get(DOWNLINK_MSGS_PERMANENTLY_FAILED.getKey()).longValue());
Assertions.assertEquals(0L, valuesByKey.get(DOWNLINK_MSGS_TMP_FAILED.getKey()).longValue());
Assertions.assertEquals(10L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey()).longValue());
Assertions.assertEquals(EXPECTED_MSGS_ADDED, valuesByKey.get(DOWNLINK_MSGS_ADDED.getKey()).longValue());
Assertions.assertEquals(EXPECTED_MSGS_PUSHED, valuesByKey.get(DOWNLINK_MSGS_PUSHED.getKey()).longValue());
Assertions.assertEquals(EXPECTED_MSGS_PERMANENTLY_FAILED, valuesByKey.get(DOWNLINK_MSGS_PERMANENTLY_FAILED.getKey()).longValue());
Assertions.assertEquals(EXPECTED_MSGS_TMP_FAILED, valuesByKey.get(DOWNLINK_MSGS_TMP_FAILED.getKey()).longValue());
return valuesByKey;
}
@Test
public void testReportStatsWithKafkaLag() {
MsgCounters counters = new MsgCounters(tenantId);
counters.getMsgsAdded().set(2);
counters.getMsgsPushed().set(2);
counters.getMsgsPermanentlyFailed().set(0);
counters.getMsgsTmpFailed().set(1);
counters.getMsgsLag().set(0);
ConcurrentHashMap<EdgeId, MsgCounters> countersByEdge = new ConcurrentHashMap<>();
countersByEdge.put(edgeId, counters);
when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge);
private void setupKafkaLag() {
String topic = "edge-topic";
TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false);
when(topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId)).thenReturn(partitionInfo);
KafkaAdmin kafkaAdmin = mock(KafkaAdmin.class);
when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic)))
.thenReturn(Map.of(topic, 15L));
when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong()))
.thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class)));
.thenReturn(Map.of(topic, EXPECTED_MSGS_KAFKA_LAG));
edgeStatsService = createEdgeStatsService(Optional.of(kafkaAdmin));
edgeStatsService.reportStats();
verify(tsService, times(1)).save(eq(tenantId), eq(edgeId), anyList(), anyLong());
verify(statsCounterService, times(1)).clear(edgeId);
List<TsKvEntry> entries = captor.getValue();
Map<String, Long> valuesByKey = entries.stream()
.collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L)));
Assertions.assertEquals(15L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey()));
}
}

3
dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java

@ -68,7 +68,8 @@ public class PostgresEdgeEventService extends BaseEdgeEventService {
Futures.addCallback(saveFuture, new FutureCallback<>() {
@Override
public void onSuccess(Void result) {
statsCounterService.ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
statsCounterService.ifPresent(statsCounterService ->
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_ADDED, edgeEvent.getTenantId(), edgeEvent.getEdgeId(), 1));
eventPublisher.publishEvent(SaveEntityEvent.builder()
.tenantId(edgeEvent.getTenantId())
.entityId(edgeEvent.getEdgeId())

Loading…
Cancel
Save