From 54e55d51ada9e90339dd81f14ed9e303ef3dfeef Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 16 Nov 2022 16:38:36 +0200 Subject: [PATCH 1/6] Use id instread of createdtime for sort order of edge events --- .../service/edge/rpc/fetch/GeneralEdgeEventFetcher.java | 2 +- .../server/dao/service/BaseEdgeEventServiceTest.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index ed5e039b62..8741426b42 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -37,7 +37,7 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher { pageSize, 0, null, - new SortOrder("createdTime", SortOrder.Direction.ASC), + new SortOrder("id", SortOrder.Direction.ASC), queueStartTs, null); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java index 73b26df9e8..11a2423ba7 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java @@ -110,7 +110,7 @@ public abstract class BaseEdgeEventServiceTest extends AbstractServiceTest { Futures.allAsList(futures).get(); - TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime); + TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("id", SortOrder.Direction.DESC), startTime, endTime); PageData edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, pageLink, true); Assert.assertNotNull(edgeEvents.getData()); @@ -135,7 +135,7 @@ public abstract class BaseEdgeEventServiceTest extends AbstractServiceTest { EdgeId edgeId = new EdgeId(Uuids.timeBased()); DeviceId deviceId = new DeviceId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); - TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime", SortOrder.Direction.ASC)); + TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("id", SortOrder.Direction.ASC)); EdgeEvent edgeEventWithTsUpdate = generateEdgeEvent(tenantId, edgeId, deviceId, EdgeEventActionType.TIMESERIES_UPDATED); edgeEventService.saveAsync(edgeEventWithTsUpdate).get(); From b5dbc7321cd33721541d3e710a5b1a42c0a9b7e6 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 18 Nov 2022 16:27:36 +0200 Subject: [PATCH 2/6] Increase timeout between batches. Improve cancel and interruption of send downlink msgs task --- .../service/edge/rpc/EdgeGrpcSession.java | 50 +++++++++++++------ .../service/edge/rpc/EdgeSessionState.java | 6 +-- .../src/main/resources/thingsboard.yml | 2 +- 3 files changed, 38 insertions(+), 20 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 8ada917234..9571319226 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 @@ -187,10 +187,7 @@ public final class EdgeGrpcSession implements Closeable { public void startSyncProcess(TenantId tenantId, EdgeId edgeId, boolean fullSync) { log.trace("[{}][{}] Staring edge sync process", tenantId, edgeId); syncCompleted = false; - if (sessionState.getSendDownlinkMsgsFuture() != null && sessionState.getSendDownlinkMsgsFuture().isDone()) { - String errorMsg = String.format("[%s][%s] Sync process started. General processing interrupted!", tenantId, edgeId); - sessionState.getSendDownlinkMsgsFuture().setException(new RuntimeException(errorMsg)); - } + interruptGeneralProcessingOnSync(tenantId, edgeId); doSync(new EdgeSyncCursor(ctx, edge, fullSync)); } @@ -265,10 +262,7 @@ public final class EdgeGrpcSession implements Closeable { } if (sessionState.getPendingMsgsMap().isEmpty()) { log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey()); - if (sessionState.getScheduledSendDownlinkTask() != null) { - sessionState.getScheduledSendDownlinkTask().cancel(false); - } - sessionState.getSendDownlinkMsgsFuture().set(null); + stopCurrentSendDownlinkMsgsTask(null); } } catch (Exception e) { log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e); @@ -391,15 +385,14 @@ public final class EdgeGrpcSession implements Closeable { } private ListenableFuture sendDownlinkMsgsPack(List downlinkMsgsPack) { - if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { - String errorMsg = "[" + this.sessionId + "] Previous send downlink future was not properly completed, stopping it now"; - log.error(errorMsg); - sessionState.getSendDownlinkMsgsFuture().setException(new RuntimeException(errorMsg)); - } + interruptPreviousSendDownlinkMsgsTask(); + sessionState.setSendDownlinkMsgsFuture(SettableFuture.create()); sessionState.getPendingMsgsMap().clear(); + downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); scheduleDownlinkMsgsPackSend(1); + return sessionState.getSendDownlinkMsgsFuture(); } @@ -422,13 +415,13 @@ public final class EdgeGrpcSession implements Closeable { } else { log.warn("[{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy); - sessionState.getSendDownlinkMsgsFuture().set(null); + stopCurrentSendDownlinkMsgsTask(null); } } else { - sessionState.getSendDownlinkMsgsFuture().set(null); + stopCurrentSendDownlinkMsgsTask(null); } } catch (Exception e) { - sessionState.getSendDownlinkMsgsFuture().setException(e); + stopCurrentSendDownlinkMsgsTask(e); } }; @@ -673,4 +666,29 @@ public final class EdgeGrpcSession implements Closeable { log.debug("[{}] Failed to close output stream: {}", sessionId, e.getMessage()); } } + + private void interruptPreviousSendDownlinkMsgsTask() { + String msg = String.format("[%s] Previous send downlink future was not properly completed, stopping it now!", this.sessionId); + stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg)); + } + + private void interruptGeneralProcessingOnSync(TenantId tenantId, EdgeId edgeId) { + String msg = String.format("[%s][%s] Sync process started. General processing interrupted!", tenantId, edgeId); + stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg)); + } + + public void stopCurrentSendDownlinkMsgsTask(Exception e) { + if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { + if (e != null) { + log.warn(e.getMessage(), e); + sessionState.getSendDownlinkMsgsFuture().setException(e); + } else { + sessionState.getSendDownlinkMsgsFuture().set(null); + } + } + if (sessionState.getScheduledSendDownlinkTask() != null) { + sessionState.getScheduledSendDownlinkTask().cancel(true); + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java index 98673c9cf5..f36ea73750 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java @@ -19,14 +19,14 @@ import com.google.common.util.concurrent.SettableFuture; import lombok.Data; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; -import java.util.LinkedHashMap; -import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledFuture; @Data public class EdgeSessionState { - private final Map pendingMsgsMap = new LinkedHashMap<>(); + private final ConcurrentMap pendingMsgsMap = new ConcurrentHashMap<>(); private SettableFuture sendDownlinkMsgsFuture; private ScheduledFuture scheduledSendDownlinkTask; } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 0f110da90c..35c613cae6 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -928,7 +928,7 @@ edges: storage: max_read_records_count: "${EDGES_STORAGE_MAX_READ_RECORDS_COUNT:50}" no_read_records_sleep: "${EDGES_NO_READ_RECORDS_SLEEP:1000}" - sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:1000}" + sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:10000}" scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}" send_scheduler_pool_size: "${EDGES_SEND_SCHEDULER_POOL_SIZE:1}" grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}" From 833ea1b3976fb56e2dd7445988ab2792fdbdf62c Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 18 Nov 2022 19:23:21 +0200 Subject: [PATCH 3/6] Decreased number of message in edge random timeseries failure to speed up tests execution --- .../java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java index 25280de666..8930166396 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java @@ -34,7 +34,7 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest { @Test public void testTimeseriesWithFailures() throws Exception { - int numberOfTimeseriesToSend = 1000; + int numberOfTimeseriesToSend = 333; Device device = findDeviceByName("Edge Device 1"); From a094b71a245059102d11e3e233f087b5413d35b2 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 21 Nov 2022 18:30:25 +0200 Subject: [PATCH 4/6] Fix stability of testSendDeviceToCloudWithNameThatAlreadyExistsOnCloud --- .../thingsboard/server/edge/BaseDeviceEdgeTest.java | 12 ++++++------ .../server/edge/imitator/EdgeImitator.java | 6 +----- 2 files changed, 7 insertions(+), 11 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java index ffe21ae961..a3162ebc5f 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java @@ -439,9 +439,9 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { Assert.assertTrue(edgeImitator.waitForResponses()); Assert.assertTrue(edgeImitator.waitForMessages()); - AbstractMessage latestMessage = edgeImitator.getMessageFromTail(2); - Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg); - DeviceUpdateMsg latestDeviceUpdateMsg = (DeviceUpdateMsg) latestMessage; + Optional deviceUpdateMsgOpt = edgeImitator.findMessageByType(DeviceUpdateMsg.class); + Assert.assertTrue(deviceUpdateMsgOpt.isPresent()); + DeviceUpdateMsg latestDeviceUpdateMsg = deviceUpdateMsgOpt.get(); Assert.assertNotEquals(deviceOnCloudName, latestDeviceUpdateMsg.getName()); Assert.assertEquals(deviceOnCloudName, latestDeviceUpdateMsg.getConflictName()); @@ -453,9 +453,9 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { Assert.assertNotNull(device); Assert.assertNotEquals(deviceOnCloudName, device.getName()); - latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof DeviceCredentialsRequestMsg); - DeviceCredentialsRequestMsg latestDeviceCredentialsRequestMsg = (DeviceCredentialsRequestMsg) latestMessage; + Optional deviceCredentialsUpdateMsgOpt = edgeImitator.findMessageByType(DeviceCredentialsRequestMsg.class); + Assert.assertTrue(deviceCredentialsUpdateMsgOpt.isPresent()); + DeviceCredentialsRequestMsg latestDeviceCredentialsRequestMsg = deviceCredentialsUpdateMsgOpt.get(); Assert.assertEquals(uuid.getMostSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdMSB()); Assert.assertEquals(uuid.getLeastSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdLSB()); diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 668f702a87..4fa3b5e723 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -365,11 +365,7 @@ public class EdgeImitator { } public AbstractMessage getLatestMessage() { - return getMessageFromTail(1); - } - - public AbstractMessage getMessageFromTail(int offset) { - return downlinkMsgs.get(downlinkMsgs.size() - offset); + return downlinkMsgs.get(downlinkMsgs.size() - 1); } public void ignoreType(Class type) { From 263605ed52741ae1a110696cca3adff0a725d313 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 21 Nov 2022 22:14:17 +0200 Subject: [PATCH 5/6] Revert Id to createdTime --- .../service/edge/rpc/fetch/GeneralEdgeEventFetcher.java | 2 +- .../server/dao/service/BaseEdgeEventServiceTest.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index 8741426b42..ed5e039b62 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -37,7 +37,7 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher { pageSize, 0, null, - new SortOrder("id", SortOrder.Direction.ASC), + new SortOrder("createdTime", SortOrder.Direction.ASC), queueStartTs, null); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java index 11a2423ba7..73b26df9e8 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeEventServiceTest.java @@ -110,7 +110,7 @@ public abstract class BaseEdgeEventServiceTest extends AbstractServiceTest { Futures.allAsList(futures).get(); - TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("id", SortOrder.Direction.DESC), startTime, endTime); + TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime); PageData edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, pageLink, true); Assert.assertNotNull(edgeEvents.getData()); @@ -135,7 +135,7 @@ public abstract class BaseEdgeEventServiceTest extends AbstractServiceTest { EdgeId edgeId = new EdgeId(Uuids.timeBased()); DeviceId deviceId = new DeviceId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); - TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("id", SortOrder.Direction.ASC)); + TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime", SortOrder.Direction.ASC)); EdgeEvent edgeEventWithTsUpdate = generateEdgeEvent(tenantId, edgeId, deviceId, EdgeEventActionType.TIMESERIES_UPDATED); edgeEventService.saveAsync(edgeEventWithTsUpdate).get(); From a0a5549b26bfb4f3dc351b01994525b3a35c60af Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 21 Nov 2022 22:56:48 +0200 Subject: [PATCH 6/6] Revert LinkedHashMap --- .../server/service/edge/rpc/EdgeSessionState.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java index f36ea73750..5a3d8dc0f8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java @@ -19,14 +19,15 @@ import com.google.common.util.concurrent.SettableFuture; import lombok.Data; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.Map; import java.util.concurrent.ScheduledFuture; @Data public class EdgeSessionState { - private final ConcurrentMap pendingMsgsMap = new ConcurrentHashMap<>(); + private final Map pendingMsgsMap = Collections.synchronizedMap(new LinkedHashMap<>()); private SettableFuture sendDownlinkMsgsFuture; private ScheduledFuture scheduledSendDownlinkTask; }