From e08f627ae125f5e6ad0ff5f4dcc8b4066552d4d1 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Thu, 14 Aug 2025 16:09:41 +0300 Subject: [PATCH 1/5] Improve Edge stats reporting and add related integration tests --- .../service/edge/rpc/EdgeGrpcSession.java | 2 +- .../service/edge/stats/EdgeStatsService.java | 18 +- .../server/edge/EdgeStatsIntegrationTest.java | 293 ++++++++++++++++++ .../server/service/edge/EdgeStatsTest.java | 68 ++-- .../thingsboard/edge/rpc/EdgeRpcClient.java | 1 + common/edge-api/src/main/proto/edge.proto | 1 + .../server/dao/edge/stats/EdgeStats.java | 34 ++ .../edge/stats/EdgeStatsCounterService.java | 16 +- 8 files changed, 381 insertions(+), 52 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 521730741f..c9364accd5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -301,7 +301,7 @@ public abstract class EdgeGrpcSession implements Closeable { if (isConnected() && !pageData.getData().isEmpty()) { if (fetcher instanceof GeneralEdgeEventFetcher) { long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage()); - ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.setDownlinkMsgsLag(edge.getTenantId(), edge.getId(), queueSize)); + ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, edge.getId(), queueSize)); } log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size()); List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java index 48b2a47cfb..697f7b6216 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.edge.stats.EdgeStats; import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService; import org.thingsboard.server.dao.edge.stats.MsgCounters; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -80,13 +81,14 @@ public class EdgeStatsService { long now = System.currentTimeMillis(); long ts = now - (now % reportIntervalMillis); - Map countersByEdge = statsCounterService.getCounterByEdge(); - Map lagByEdgeId = kafkaAdmin.isPresent() ? getEdgeLagByEdgeId(countersByEdge) : Collections.emptyMap(); - Map countersByEdgeSnapshot = new HashMap<>(statsCounterService.getCounterByEdge()); - countersByEdgeSnapshot.forEach((edgeId, counters) -> { + Map statsByEdgeSnapshot = new HashMap<>(statsCounterService.getStatsByEdge()); + boolean isKafkaStats = kafkaAdmin.isPresent(); + Map lagByEdgeId = isKafkaStats ? getLagByEdgeId(statsByEdgeSnapshot) : Collections.emptyMap(); + statsByEdgeSnapshot.forEach((edgeId, edgeStats) -> { + MsgCounters counters = edgeStats.getMsgCounters(); TenantId tenantId = counters.getTenantId(); - if (kafkaAdmin.isPresent()) { + if (isKafkaStats) { counters.getMsgsLag().set(lagByEdgeId.getOrDefault(edgeId, 0L)); } List statsEntries = List.of( @@ -102,11 +104,11 @@ public class EdgeStatsService { }); } - private Map getEdgeLagByEdgeId(Map countersByEdge) { - Map edgeToTopicMap = countersByEdge.entrySet().stream() + private Map getLagByEdgeId(Map edgeStatsByEdge) { + Map edgeToTopicMap = edgeStatsByEdge.entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, - e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getTenantId(), e.getKey()).getTopic() + e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getMsgCounters().getTenantId(), e.getKey()).getTopic() )); Map lagByTopic = kafkaAdmin.get().getTotalLagForGroupsBulk(new HashSet<>(edgeToTopicMap.values())); diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java new file mode 100644 index 0000000000..b993e4f564 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java @@ -0,0 +1,293 @@ +/** + * 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.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; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.kv.TsKvEntry; +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; +import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_TMP_FAILED; + +@DaoSqlTest +@Slf4j +public class EdgeStatsIntegrationTest extends AbstractEdgeTest { + + private static final String STATISTICS_DEVICE_PROFILE = "STATISTICS"; + + private final Map 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 EXPECTED_EMPTY_EDGE_STATS = Map.of( + DOWNLINK_MSGS_ADDED, 0L, + DOWNLINK_MSGS_PUSHED, 0L, + DOWNLINK_MSGS_PERMANENTLY_FAILED, 0L, + DOWNLINK_MSGS_TMP_FAILED, 0L + ); + + @Autowired + private EdgeStatsService edgeStatsService; + @Autowired + private EdgeStatsCounterService statsCounterService; + + @Test + public void testFullEdgeStatsCycle() throws Exception { + // 1. Clear previous stats and prepare test data + 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 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().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> { + EntityDataProto latestMsg = getLatestEntityDataMessage(); + Map actualStats = toMap(latestMsg.getPostTelemetryMsg().getTsKvList(0)); + assertAllStatsEqual(EXPECTED_EDGE_STATS, actualStats, "Proto stats"); + }); + } + + @Test + public void testNoMessagesFromEdge() throws ExecutionException, InterruptedException { + // 1. Clear stats counters for the Edge + statsCounterService.clear(edge.getId()); + + // 2. Report stats with no data from the Edge + edgeStatsService.reportStats(); + + // 3. Verify that persisted timeseries contains only empty stats + awaitStatsMatch(EXPECTED_EMPTY_EDGE_STATS); + + List 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 expected) { + await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> { + Map actualStats = fetchLatestStats(); + assertAllStatsEqual(expected, actualStats, "Timeseries stats"); + }); + } + + private Map fetchLatestStats() throws ExecutionException, InterruptedException { + List latestStatsEntries = tsService.findLatest( + tenantId, + edge.getId(), + Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList() + ).get(); + return toMap(latestStatsEntries); + } + + private void prepareTestData() throws InterruptedException, ExecutionException { + statsCounterService.clear(edge.getId()); + // 2 stats message ADDED Device Profile, ASSIGN Device + 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 + edgeImitator.expectMessageAmount(2); + Asset savedAsset = saveAsset("Edge Asset 2"); + doPost("/api/edge/" + edge.getUuidId() + + "/asset/" + savedAsset.getUuidId(), Asset.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + // 2 stats message ADDED Customer, ASSIGN Customer + edgeImitator.expectMessageAmount(1); + 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 + 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); + edgeEventService.saveAsync(edgeEvent).get(); + Assert.assertTrue(edgeImitator.waitForMessages()); + } + + private void assertAllStatsEqual(Map expected, Map 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.getStatsByEdge().get(edge.getId()).getMsgCounters(); + Map 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 final Map toMap(Map.Entry>... suppliers) { + return Arrays.stream(suppliers) + .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get())); + } + + private Map toMap(List stats) { + return stats.stream() + .collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(0L))); + } + + private Map toMap(TransportProtos.TsKvListProto kvList) { + Map 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 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 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; + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java index 25ff0f1b5d..fb24cdb5bb 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.dao.edge.stats.EdgeStats; import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService; import org.thingsboard.server.dao.edge.stats.MsgCounters; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -44,9 +45,11 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; +import static org.mockito.ArgumentMatchers.anyList; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_ADDED; @@ -58,6 +61,9 @@ import static org.thingsboard.server.dao.edge.stats.EdgeStatsKey.DOWNLINK_MSGS_T @ExtendWith(MockitoExtension.class) public class EdgeStatsTest { + private static final int TTL_DAYS = 30; + private static final long REPORT_INTERVAL_MILLIS = 600_000L; + @Mock private TimeseriesService tsService; @Mock @@ -71,40 +77,45 @@ public class EdgeStatsTest { @BeforeEach void setUp() { - edgeStatsService = new EdgeStatsService( + edgeStatsService = createEdgeStatsService(Optional.empty()); + } + + private EdgeStatsService createEdgeStatsService(Optional kafkaAdmin) { + EdgeStatsService service = new EdgeStatsService( tsService, statsCounterService, topicService, - Optional.empty() + kafkaAdmin ); - - ReflectionTestUtils.setField(edgeStatsService, "edgesStatsTtlDays", 30); - ReflectionTestUtils.setField(edgeStatsService, "reportIntervalMillis", 600_000L); + ReflectionTestUtils.setField(service, "edgesStatsTtlDays", TTL_DAYS); + ReflectionTestUtils.setField(service, "reportIntervalMillis", REPORT_INTERVAL_MILLIS); + return service; } @Test public void testReportStatsSavesTelemetry() { - // given - MsgCounters counters = new MsgCounters(tenantId); + EdgeStats edgeStats = new EdgeStats(tenantId); + MsgCounters counters = edgeStats.getMsgCounters(); counters.getMsgsAdded().set(5); counters.getMsgsPushed().set(3); counters.getMsgsPermanentlyFailed().set(1); counters.getMsgsTmpFailed().set(0); counters.getMsgsLag().set(10); - ConcurrentHashMap countersByEdge = new ConcurrentHashMap<>(); - countersByEdge.put(edgeId, counters); + ConcurrentHashMap edgeStatsByEdge = new ConcurrentHashMap<>(); + edgeStatsByEdge.put(edgeId, edgeStats); - when(statsCounterService.getCounterByEdge()).thenReturn(countersByEdge); + when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge); - ArgumentCaptor> captor = ArgumentCaptor.forClass(List.class); + ArgumentCaptor> captor = ArgumentCaptor.forClass((Class) List.class); when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); - // when edgeStatsService.reportStats(); - // then + verify(tsService, times(1)).save(eq(tenantId), eq(edgeId), anyList(), anyLong()); + verify(statsCounterService, times(1)).clear(edgeId); + List entries = captor.getValue(); Assertions.assertEquals(5, entries.size()); @@ -116,26 +127,22 @@ public class EdgeStatsTest { 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()); - - - verify(statsCounterService).clear(edgeId); } @Test public void testReportStatsWithKafkaLag() { - // given - MsgCounters counters = new MsgCounters(tenantId); + EdgeStats edgeStats = new EdgeStats(tenantId); + MsgCounters counters = edgeStats.getMsgCounters(); counters.getMsgsAdded().set(2); counters.getMsgsPushed().set(2); counters.getMsgsPermanentlyFailed().set(0); counters.getMsgsTmpFailed().set(1); counters.getMsgsLag().set(0); - ConcurrentHashMap countersByEdge = new ConcurrentHashMap<>(); - countersByEdge.put(edgeId, counters); + ConcurrentHashMap edgeStatsByEdge = new ConcurrentHashMap<>(); + edgeStatsByEdge.put(edgeId, edgeStats); - // mocks - when(statsCounterService.getCounterByEdge()).thenReturn(countersByEdge); + when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge); String topic = "edge-topic"; TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false); @@ -145,29 +152,22 @@ public class EdgeStatsTest { when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic))) .thenReturn(Map.of(topic, 15L)); - ArgumentCaptor> captor = ArgumentCaptor.forClass(List.class); + ArgumentCaptor> captor = ArgumentCaptor.forClass((Class) List.class); when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); - edgeStatsService = new EdgeStatsService( - tsService, - statsCounterService, - topicService, - Optional.of(kafkaAdmin) - ); - ReflectionTestUtils.setField(edgeStatsService, "edgesStatsTtlDays", 30); - ReflectionTestUtils.setField(edgeStatsService, "reportIntervalMillis", 600_000L); + edgeStatsService = createEdgeStatsService(Optional.of(kafkaAdmin)); - // when edgeStatsService.reportStats(); - // then + verify(tsService, times(1)).save(eq(tenantId), eq(edgeId), anyList(), anyLong()); + verify(statsCounterService, times(1)).clear(edgeId); + List entries = captor.getValue(); Map valuesByKey = entries.stream() .collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L))); Assertions.assertEquals(15L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey())); - verify(statsCounterService).clear(edgeId); } } diff --git a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java index 423c59251b..ed0d5b603e 100644 --- a/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java +++ b/common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeRpcClient.java @@ -41,4 +41,5 @@ public interface EdgeRpcClient { void sendDownlinkResponseMsg(DownlinkResponseMsg downlinkResponseMsg); int getServerMaxInboundMessageSize(); + } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index dbda462a99..aea505c2fe 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -44,6 +44,7 @@ enum EdgeVersion { V_4_0_0 = 10; V_4_1_0 = 11; V_4_2_0 = 12; + V_4_3_0 = 13; V_LATEST = 999; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java new file mode 100644 index 0000000000..6f34563f96 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java @@ -0,0 +1,34 @@ +/** + * 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.dao.edge.stats; + +import lombok.Data; +import org.thingsboard.server.common.data.id.TenantId; + +import java.util.Queue; +import java.util.concurrent.ConcurrentLinkedQueue; + +@Data +public class EdgeStats { + private final MsgCounters msgCounters; + private final Queue uplinkRate; + + public EdgeStats(TenantId tenantId) { + this.msgCounters = new MsgCounters(tenantId); + this.uplinkRate = new ConcurrentLinkedQueue<>(); + } + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java index 16111cf514..c537fb75be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java @@ -30,28 +30,26 @@ import java.util.concurrent.ConcurrentHashMap; @Getter public class EdgeStatsCounterService { - private final ConcurrentHashMap counterByEdge = new ConcurrentHashMap<>(); + private final ConcurrentHashMap statsByEdge = new ConcurrentHashMap<>(); public void recordEvent(EdgeStatsKey type, TenantId tenantId, EdgeId edgeId, long value) { - MsgCounters counters = getOrCreateCounters(tenantId, edgeId); + EdgeStats edgeStats = getOrCreateEdgeStats(tenantId, edgeId); + MsgCounters counters = edgeStats.getMsgCounters(); switch (type) { case DOWNLINK_MSGS_ADDED -> counters.getMsgsAdded().addAndGet(value); case DOWNLINK_MSGS_PUSHED -> counters.getMsgsPushed().addAndGet(value); case DOWNLINK_MSGS_PERMANENTLY_FAILED -> counters.getMsgsPermanentlyFailed().addAndGet(value); case DOWNLINK_MSGS_TMP_FAILED -> counters.getMsgsTmpFailed().addAndGet(value); + case DOWNLINK_MSGS_LAG -> counters.getMsgsLag().set(value); } } - public void setDownlinkMsgsLag(TenantId tenantId, EdgeId edgeId, long value) { - getOrCreateCounters(tenantId, edgeId).getMsgsLag().set(value); + public EdgeStats getOrCreateEdgeStats(TenantId tenantId, EdgeId edgeId) { + return statsByEdge.computeIfAbsent(edgeId, id -> new EdgeStats(tenantId)); } public void clear(EdgeId edgeId) { - counterByEdge.remove(edgeId); - } - - public MsgCounters getOrCreateCounters(TenantId tenantId, EdgeId edgeId) { - return counterByEdge.computeIfAbsent(edgeId, id -> new MsgCounters(tenantId)); + statsByEdge.remove(edgeId); } } From d35be826b8f68239da443c116c30bdba5732d497 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Thu, 14 Aug 2025 16:13:10 +0300 Subject: [PATCH 2/5] Refactor --- .../java/org/thingsboard/server/dao/edge/stats/EdgeStats.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java index 6f34563f96..36d87691cc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java @@ -31,4 +31,4 @@ public class EdgeStats { this.uplinkRate = new ConcurrentLinkedQueue<>(); } -} \ No newline at end of file +} From e8b164a03ad2e07fa3c8ecc61daf14e70d0ed7dc Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 9 Dec 2025 16:56:32 +0200 Subject: [PATCH 3/5] Refactoring. Remove EdgeStats - uplinkRate not used --- .../service/edge/EdgeContextComponent.java | 1 - .../service/edge/rpc/EdgeGrpcSession.java | 6 ++-- .../service/edge/stats/EdgeStatsService.java | 14 ++++---- .../server/edge/EdgeStatsIntegrationTest.java | 12 +++---- .../server/service/edge/EdgeStatsTest.java | 25 +++++++------- .../server/dao/edge/stats/EdgeStats.java | 34 ------------------- .../edge/stats/EdgeStatsCounterService.java | 13 ++++--- 7 files changed, 34 insertions(+), 71 deletions(-) delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index ec4a43310b..d523a9176c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -206,7 +206,6 @@ public class EdgeContextComponent { @Autowired private Optional statsCounterService; - // processors @Autowired private AlarmProcessor alarmProcessor; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index c0e76c32e4..83639573bb 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -306,7 +306,8 @@ public abstract class EdgeGrpcSession implements Closeable { if (isConnected() && !pageData.getData().isEmpty()) { if (fetcher instanceof GeneralEdgeEventFetcher) { long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage()); - ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, edge.getId(), queueSize)); + ctx.getStatsCounterService().ifPresent(statsCounterService -> + statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, edge.getId(), queueSize)); } log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size()); List downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); @@ -504,7 +505,8 @@ public abstract class EdgeGrpcSession implements Closeable { ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) .customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) .error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); - ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), edge.getId(), copy.size())); + ctx.getStatsCounterService().ifPresent(statsCounterService -> + statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), edge.getId(), copy.size())); stopCurrentSendDownlinkMsgsTask(false); } } else { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java index 697f7b6216..66be8fe32a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java @@ -31,7 +31,6 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; import org.thingsboard.server.common.data.kv.TsKvEntry; -import org.thingsboard.server.dao.edge.stats.EdgeStats; import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService; import org.thingsboard.server.dao.edge.stats.MsgCounters; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -81,11 +80,10 @@ public class EdgeStatsService { long now = System.currentTimeMillis(); long ts = now - (now % reportIntervalMillis); - Map statsByEdgeSnapshot = new HashMap<>(statsCounterService.getStatsByEdge()); + Map countersByEdgeSnapshot = new HashMap<>(statsCounterService.getMsgCountersByEdge()); boolean isKafkaStats = kafkaAdmin.isPresent(); - Map lagByEdgeId = isKafkaStats ? getLagByEdgeId(statsByEdgeSnapshot) : Collections.emptyMap(); - statsByEdgeSnapshot.forEach((edgeId, edgeStats) -> { - MsgCounters counters = edgeStats.getMsgCounters(); + Map lagByEdgeId = isKafkaStats ? getLagByEdgeId(countersByEdgeSnapshot) : Collections.emptyMap(); + countersByEdgeSnapshot.forEach((edgeId, counters) -> { TenantId tenantId = counters.getTenantId(); if (isKafkaStats) { @@ -104,11 +102,11 @@ public class EdgeStatsService { }); } - private Map getLagByEdgeId(Map edgeStatsByEdge) { - Map edgeToTopicMap = edgeStatsByEdge.entrySet().stream() + private Map getLagByEdgeId(Map countersByEdge) { + Map edgeToTopicMap = countersByEdge.entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, - e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getMsgCounters().getTenantId(), e.getKey()).getTopic() + e -> topicService.buildEdgeEventNotificationsTopicPartitionInfo(e.getValue().getTenantId(), e.getKey()).getTopic() )); Map lagByTopic = kafkaAdmin.get().getTotalLagForGroupsBulk(new HashSet<>(edgeToTopicMap.values())); diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java index b993e4f564..ed98ebe16d 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java @@ -179,19 +179,19 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { private void prepareTestData() throws InterruptedException, ExecutionException { statsCounterService.clear(edge.getId()); - // 2 stats message ADDED Device Profile, ASSIGN Device + // 2 stats messages: ADDED Device Profile, ASSIGN Device 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 + // 1 stats message: ASSIGN Asset edgeImitator.expectMessageAmount(2); Asset savedAsset = saveAsset("Edge Asset 2"); doPost("/api/edge/" + edge.getUuidId() + "/asset/" + savedAsset.getUuidId(), Asset.class); Assert.assertTrue(edgeImitator.waitForMessages()); - // 2 stats message ADDED Customer, ASSIGN Customer + // 2 stats messages: ADDED Customer, ASSIGN Customer edgeImitator.expectMessageAmount(1); Customer customer = new Customer(); customer.setTitle("Edge Customer"); @@ -204,7 +204,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { + "/edge/" + edge.getUuidId(), Edge.class); Assert.assertTrue(edgeImitator.waitForMessages()); - //1 stats message Timeseries + // 1 stats message: Timeseries edgeImitator.expectMessageAmount(1); String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}"; JsonNode timeseriesEntityData = JacksonUtil.toJsonNode(timeseriesData); @@ -222,7 +222,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { private void awaitEdgeCountersUpdated() { await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> { - MsgCounters counters = statsCounterService.getStatsByEdge().get(edge.getId()).getMsgCounters(); + MsgCounters counters = statsCounterService.getMsgCountersByEdge().get(edge.getId()); Map actualCounters = toMap( Map.entry(DOWNLINK_MSGS_ADDED.getKey(), () -> counters.getMsgsAdded().get()), Map.entry(DOWNLINK_MSGS_PUSHED.getKey(), () -> counters.getMsgsPushed().get()), @@ -234,7 +234,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { } @SafeVarargs - private final Map toMap(Map.Entry>... suppliers) { + private Map toMap(Map.Entry>... suppliers) { return Arrays.stream(suppliers) .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get())); } diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java index fb24cdb5bb..0b0f28e923 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java @@ -21,6 +21,7 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.ArgumentCaptor; +import org.mockito.Captor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.test.util.ReflectionTestUtils; @@ -29,7 +30,6 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.dao.edge.stats.EdgeStats; import org.thingsboard.server.dao.edge.stats.EdgeStatsCounterService; import org.thingsboard.server.dao.edge.stats.MsgCounters; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -72,6 +72,9 @@ public class EdgeStatsTest { private EdgeStatsCounterService statsCounterService; private EdgeStatsService edgeStatsService; + @Captor + private ArgumentCaptor> captor; + private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); private final EdgeId edgeId = new EdgeId(UUID.randomUUID()); @@ -94,20 +97,18 @@ public class EdgeStatsTest { @Test public void testReportStatsSavesTelemetry() { - EdgeStats edgeStats = new EdgeStats(tenantId); - MsgCounters counters = edgeStats.getMsgCounters(); + 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); - ConcurrentHashMap edgeStatsByEdge = new ConcurrentHashMap<>(); - edgeStatsByEdge.put(edgeId, edgeStats); + ConcurrentHashMap countersByEdge = new ConcurrentHashMap<>(); + countersByEdge.put(edgeId, counters); - when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge); + when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge); - ArgumentCaptor> captor = ArgumentCaptor.forClass((Class) List.class); when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); @@ -131,18 +132,17 @@ public class EdgeStatsTest { @Test public void testReportStatsWithKafkaLag() { - EdgeStats edgeStats = new EdgeStats(tenantId); - MsgCounters counters = edgeStats.getMsgCounters(); + 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 edgeStatsByEdge = new ConcurrentHashMap<>(); - edgeStatsByEdge.put(edgeId, edgeStats); + ConcurrentHashMap countersByEdge = new ConcurrentHashMap<>(); + countersByEdge.put(edgeId, counters); - when(statsCounterService.getStatsByEdge()).thenReturn(edgeStatsByEdge); + when(statsCounterService.getMsgCountersByEdge()).thenReturn(countersByEdge); String topic = "edge-topic"; TopicPartitionInfo partitionInfo = new TopicPartitionInfo(topic, tenantId, 0, false); @@ -152,7 +152,6 @@ public class EdgeStatsTest { when(kafkaAdmin.getTotalLagForGroupsBulk(Set.of(topic))) .thenReturn(Map.of(topic, 15L)); - ArgumentCaptor> captor = ArgumentCaptor.forClass((Class) List.class); when(tsService.save(eq(tenantId), eq(edgeId), captor.capture(), anyLong())) .thenReturn(Futures.immediateFuture(mock(TimeseriesSaveResult.class))); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java deleted file mode 100644 index 36d87691cc..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStats.java +++ /dev/null @@ -1,34 +0,0 @@ -/** - * 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.dao.edge.stats; - -import lombok.Data; -import org.thingsboard.server.common.data.id.TenantId; - -import java.util.Queue; -import java.util.concurrent.ConcurrentLinkedQueue; - -@Data -public class EdgeStats { - private final MsgCounters msgCounters; - private final Queue uplinkRate; - - public EdgeStats(TenantId tenantId) { - this.msgCounters = new MsgCounters(tenantId); - this.uplinkRate = new ConcurrentLinkedQueue<>(); - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java index c537fb75be..ea8b85a287 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/stats/EdgeStatsCounterService.java @@ -24,17 +24,16 @@ import org.thingsboard.server.common.data.id.TenantId; import java.util.concurrent.ConcurrentHashMap; -@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true", matchIfMissing = false) +@ConditionalOnProperty(prefix = "edges.stats", name = "enabled", havingValue = "true") @Service @Slf4j @Getter public class EdgeStatsCounterService { - private final ConcurrentHashMap statsByEdge = new ConcurrentHashMap<>(); + private final ConcurrentHashMap msgCountersByEdge = new ConcurrentHashMap<>(); public void recordEvent(EdgeStatsKey type, TenantId tenantId, EdgeId edgeId, long value) { - EdgeStats edgeStats = getOrCreateEdgeStats(tenantId, edgeId); - MsgCounters counters = edgeStats.getMsgCounters(); + MsgCounters counters = getOrCreateCounters(tenantId, edgeId); switch (type) { case DOWNLINK_MSGS_ADDED -> counters.getMsgsAdded().addAndGet(value); case DOWNLINK_MSGS_PUSHED -> counters.getMsgsPushed().addAndGet(value); @@ -44,12 +43,12 @@ public class EdgeStatsCounterService { } } - public EdgeStats getOrCreateEdgeStats(TenantId tenantId, EdgeId edgeId) { - return statsByEdge.computeIfAbsent(edgeId, id -> new EdgeStats(tenantId)); + public MsgCounters getOrCreateCounters(TenantId tenantId, EdgeId edgeId) { + return msgCountersByEdge.computeIfAbsent(edgeId, id -> new MsgCounters(tenantId)); } public void clear(EdgeId edgeId) { - statsByEdge.remove(edgeId); + msgCountersByEdge.remove(edgeId); } } From c42705c03271be54c04893c57997518a60c72aad Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 10 Dec 2025 10:04:26 +0200 Subject: [PATCH 4/5] Simplify tests --- .../service/edge/rpc/EdgeGrpcSession.java | 6 +- .../edge/rpc/KafkaEdgeEventService.java | 3 +- .../service/edge/stats/EdgeStatsService.java | 3 +- .../server/edge/EdgeStatsIntegrationTest.java | 237 ++++-------------- .../server/service/edge/EdgeStatsTest.java | 88 ++++--- .../dao/edge/PostgresEdgeEventService.java | 3 +- 6 files changed, 102 insertions(+), 238 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 83639573bb..668c698376 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/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> processUplinkMsg(UplinkMsg uplinkMsg) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java index bc00ef4481..b4f7a7574f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java +++ b/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); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java b/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java index 66be8fe32a..7f55d2817c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/stats/EdgeStatsService.java +++ b/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}" diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java index ed98ebe16d..2789b0d7b8 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java +++ b/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 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 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 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 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 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 expected) { + // THEN await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> { - Map actualStats = fetchLatestStats(); - assertAllStatsEqual(expected, actualStats, "Timeseries stats"); + List 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 fetchLatestStats() throws ExecutionException, InterruptedException { - List latestStatsEntries = tsService.findLatest( + private long getStatsLongValue(List stats, EdgeStatsKey key) { + return stats.stream().filter(e -> e.getKey().equals(key.getKey())).findFirst().get().getLongValue().orElse(0L); + } + + private List 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 expected, Map 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 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 toMap(Map.Entry>... suppliers) { - return Arrays.stream(suppliers) - .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().get())); - } - - private Map toMap(List stats) { - return stats.stream() - .collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(0L))); - } - - private Map toMap(TransportProtos.TsKvListProto kvList) { - Map 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 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 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; - } - } diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java index 0b0f28e923..3cb4bec803 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeStatsTest.java +++ b/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 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 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 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 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 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 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 entries = captor.getValue(); - Map valuesByKey = entries.stream() - .collect(Collectors.toMap(TsKvEntry::getKey, e -> e.getLongValue().orElse(-1L))); - - Assertions.assertEquals(15L, valuesByKey.get(DOWNLINK_MSGS_LAG.getKey())); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java index ecf1376a73..e38b115e88 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java +++ b/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()) From 916132d6b4cbb0180df446888b75fa0b103208e0 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 10 Dec 2025 10:07:16 +0200 Subject: [PATCH 5/5] Test naming conventions --- .../thingsboard/server/edge/EdgeStatsIntegrationTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java index 2789b0d7b8..71722ad2a8 100644 --- a/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/EdgeStatsIntegrationTest.java @@ -65,9 +65,9 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { private EdgeStatsCounterService statsCounterService; @Test - public void testFullEdgeStatsCycle() throws Exception { + public void testReportStats() throws Exception { // GIVEN - prepareTestData(); + simulateEdgeEventsAddedDownlinkPushed(); // Await Edge Counters Updated await().atMost(10, TimeUnit.SECONDS).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> { @@ -104,7 +104,7 @@ public class EdgeStatsIntegrationTest extends AbstractEdgeTest { Arrays.stream(EdgeStatsKey.values()).map(EdgeStatsKey::getKey).toList()).get(); } - private void prepareTestData() throws InterruptedException, ExecutionException { + private void simulateEdgeEventsAddedDownlinkPushed() throws InterruptedException, ExecutionException { statsCounterService.clear(edge.getId()); // Save device and assign to edge