Browse Source

Merge pull request #13075 from jekka001/fix-sorting-edge-event

Fix sorting edge event
pull/13429/head
yevhenii_zahrebelnyi 1 year ago
committed by GitHub
parent
commit
4b2542bd47
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java
  3. 48
      application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java
  4. 21
      application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
  5. 12
      application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java
  6. 9
      application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java
  7. 15
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  8. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java
  9. 25
      dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java

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

@ -712,7 +712,7 @@ public abstract class EdgeGrpcSession implements Closeable {
private long findStartSeqIdFromOldestEventIfAny() {
long startSeqId = 0L;
try {
TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime"), null, null);
TimePageLink pageLink = new TimePageLink(1, 0, null, null, null, null);
PageData<EdgeEvent> edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), null, null, pageLink);
if (!edgeEvents.getData().isEmpty()) {
startSeqId = edgeEvents.getData().get(0).getSeqId() - 1;

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

@ -26,9 +26,13 @@ import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.edge.EdgeEventService;
import java.util.concurrent.TimeUnit;
@AllArgsConstructor
@Slf4j
public class GeneralEdgeEventFetcher implements EdgeEventFetcher {
// Subtract from queueStartTs to ensure no data is lost due to potential misordering of edge events by created_time.
private static final long MISORDERING_COMPENSATION_MILLIS = TimeUnit.SECONDS.toMillis(60);
private final Long queueStartTs;
private Long seqIdStart;
@ -44,7 +48,7 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher {
0,
null,
null,
queueStartTs,
queueStartTs > 0 ? queueStartTs - MISORDERING_COMPENSATION_MILLIS : 0,
System.currentTimeMillis());
}

48
application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java

@ -60,6 +60,7 @@ import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.page.TimePageLink;
@ -68,6 +69,7 @@ import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainMetaData;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.UserCredentials;
import org.thingsboard.server.common.data.security.model.JwtSettings;
import org.thingsboard.server.dao.edge.EdgeDao;
import org.thingsboard.server.dao.exception.DataValidationException;
@ -107,6 +109,7 @@ import java.util.concurrent.TimeUnit;
import static org.hamcrest.Matchers.containsString;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
import static org.thingsboard.server.edge.AbstractEdgeTest.CONNECT_MESSAGE_COUNT;
@TestPropertySource(properties = {
"edges.enabled=true",
@ -138,6 +141,7 @@ public class EdgeControllerTest extends AbstractControllerTest {
public EdgeDao edgeDao(EdgeDao edgeDao) {
return Mockito.mock(EdgeDao.class, AdditionalAnswers.delegatesTo(edgeDao));
}
}
@Before
@ -886,6 +890,8 @@ public class EdgeControllerTest extends AbstractControllerTest {
Device savedDevice = doPost("/api/device", device, Device.class);
// create public customer
//1 message
// Customer
doPost("/api/customer/public/device/" + savedDevice.getId().getId(), Device.class);
doDelete("/api/customer/device/" + savedDevice.getId().getId(), Device.class);
@ -897,13 +903,16 @@ public class EdgeControllerTest extends AbstractControllerTest {
+ "/asset/" + savedAsset.getId().getId().toString(), Asset.class);
EdgeImitator edgeImitator = new EdgeImitator(EDGE_HOST, EDGE_PORT, edge.getRoutingKey(), edge.getSecret());
edgeImitator.ignoreType(UserCredentialsUpdateMsg.class);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.expectMessageAmount(27);
// 17 connect message
// + 1 Customer
// + 5 fetchers messages (DeviceProfile, Device, DeviceCredentials, AssetProfile, Asset) in sync process
// + 5 queue messages the same
edgeImitator.expectMessageAmount(CONNECT_MESSAGE_COUNT + 11);
edgeImitator.connect();
waitForMessages(edgeImitator);
edgeImitator.waitForMessages();
verifyFetchersMsgs(edgeImitator, savedDevice);
// verify queue msgs
@ -914,9 +923,12 @@ public class EdgeControllerTest extends AbstractControllerTest {
Assert.assertTrue(popAssetMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Test Sync Edge Asset 1"));
printQueueMsgsIfNotEmpty(edgeImitator);
edgeImitator.expectMessageAmount(21);
// 17 connect messages
// + 1 Customer
// + 5 fetchers messages (DeviceProfile, Device, DeviceCredentials, AssetProfile, Asset) in sync process
edgeImitator.expectMessageAmount(CONNECT_MESSAGE_COUNT + 6);
doPost("/api/edge/sync/" + edge.getId()).andExpect(status().isOk());
waitForMessages(edgeImitator);
edgeImitator.waitForMessages();
verifyFetchersMsgs(edgeImitator, savedDevice);
printQueueMsgsIfNotEmpty(edgeImitator);
@ -987,17 +999,6 @@ public class EdgeControllerTest extends AbstractControllerTest {
});
}
private void waitForMessages(EdgeImitator edgeImitator) throws Exception {
boolean success = edgeImitator.waitForMessages();
if (!success) {
List<AbstractMessage> downlinkMsgs = edgeImitator.getDownlinkMsgs();
for (AbstractMessage downlinkMsg : downlinkMsgs) {
log.error("{}\n{}", downlinkMsg.getClass(), downlinkMsg);
}
Assert.fail("Await for messages was not successful!");
}
}
private void verifyFetchersMsgs(EdgeImitator edgeImitator, Device savedDevice) {
Assert.assertTrue(popQueueMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Main"));
Assert.assertTrue(popRuleChainMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Edge Root Rule Chain"));
@ -1011,6 +1012,7 @@ public class EdgeControllerTest extends AbstractControllerTest {
Assert.assertTrue(popAssetProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default"));
Assert.assertTrue(popDeviceProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default"));
Assert.assertTrue(popAssetProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default"));
Assert.assertTrue(popUserCredentialsMsg(edgeImitator.getDownlinkMsgs(), currentUserId));
Assert.assertTrue(popUserMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, TENANT_ADMIN_EMAIL, Authority.TENANT_ADMIN));
Assert.assertTrue(popCustomerMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Public"));
Assert.assertTrue(popDeviceProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default"));
@ -1156,6 +1158,20 @@ public class EdgeControllerTest extends AbstractControllerTest {
return false;
}
private boolean popUserCredentialsMsg(List<AbstractMessage> messages, UserId userId) {
for (AbstractMessage message : messages) {
if (message instanceof UserCredentialsUpdateMsg userCredentialsUpdateMsg) {
UserCredentials userCredentials = JacksonUtil.fromString(userCredentialsUpdateMsg.getEntity(), UserCredentials.class, true);
Assert.assertNotNull(userCredentials);
if (userId.equals(userCredentials.getUserId())) {
messages.remove(message);
return true;
}
}
}
return false;
}
private boolean popUserMsg(List<AbstractMessage> messages, UpdateMsgType msgType, String email, Authority authority) {
for (AbstractMessage message : messages) {
if (message instanceof UserUpdateMsg userUpdateMsg) {

21
application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java

@ -115,7 +115,9 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
})
@Slf4j
abstract public class AbstractEdgeTest extends AbstractControllerTest {
public static final Integer CONNECT_MESSAGE_COUNT = 17;
public static final Integer INSTALLATION_MESSAGE_COUNT = 8;
public static final Integer SYNC_MESSAGE_COUNT = CONNECT_MESSAGE_COUNT + INSTALLATION_MESSAGE_COUNT;
private static final String THERMOSTAT_DEVICE_PROFILE_NAME = "Thermostat";
protected DeviceProfile thermostatDeviceProfile;
@ -136,11 +138,12 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
doPost("/api/admin/jwtSettings", settings).andExpect(status().isOk());
loginTenantAdmin();
//8 installation messages
installation();
edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret());
edgeImitator.expectMessageAmount(25);
// 17 connect messages + 8 installation messages
edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.connect();
@ -164,22 +167,32 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
thermostatDeviceProfile = this.createDeviceProfile(THERMOSTAT_DEVICE_PROFILE_NAME,
createMqttDeviceProfileTransportConfiguration(new JsonTransportPayloadConfiguration(), false));
extendDeviceProfileData(thermostatDeviceProfile);
//2 messages DeviceProfile
thermostatDeviceProfile = doPost("/api/deviceProfile", thermostatDeviceProfile, DeviceProfile.class);
Device savedDevice = saveDevice("Edge Device 1", THERMOSTAT_DEVICE_PROFILE_NAME);
// create public customer
//1 message
// Customer
doPost("/api/customer/public/device/" + savedDevice.getId().getId(), Device.class);
doDelete("/api/customer/device/" + savedDevice.getId().getId(), Device.class);
Asset savedAsset = saveAsset("Edge Asset 1");
Asset savedAsset = saveAsset("Edge Asset 1");
updateRootRuleChainMetadata();
edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class);
//3 messages
// Device
// DeviceProfile
// DeviceCredentials
doPost("/api/edge/" + edge.getUuidId()
+ "/device/" + savedDevice.getUuidId(), Device.class);
//2 messages
// Asset
// AssetProfile
doPost("/api/edge/" + edge.getUuidId()
+ "/asset/" + savedAsset.getUuidId(), Asset.class);

12
application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java

@ -128,15 +128,19 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
@Test
public void testDeleteDeviceProfilesWhenEdgeIsOffline() throws Exception {
//2 message RuleChain and RuleChainMetadata
RuleChainId thermostatsRuleChainId = createEdgeRuleChainAndAssignToEdge("Thermostats Rule Chain");
// create device profile
DeviceProfile deviceProfile = this.createDeviceProfile("ONE_MORE_DEVICE_PROFILE", null);
deviceProfile.setDefaultEdgeRuleChainId(thermostatsRuleChainId);
extendDeviceProfileData(deviceProfile);
//1 message DeviceProfile
edgeImitator.expectMessageAmount(1);
deviceProfile = doPost("/api/deviceProfile", deviceProfile, DeviceProfile.class);
Assert.assertTrue(edgeImitator.waitForMessages());
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof DeviceProfileUpdateMsg);
DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage;
@ -150,9 +154,11 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
doDelete("/api/deviceProfile/" + deviceProfile.getUuidId())
.andExpect(status().isOk());
edgeImitator.connect();
// 27 sync message
// + 1 delete message
edgeImitator.expectMessageAmount(28);
// 25 sync message
// + 2 RuleChain and RuleChainMetadata
// + 1 delete DeviceProfile
edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 3);
Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage();

9
application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java

@ -78,11 +78,14 @@ public class TenantProfileEdgeTest extends AbstractEdgeTest {
TenantProfileQueueConfiguration mainQueueConfiguration = createQueueConfig(DataConstants.MAIN_QUEUE_NAME, DataConstants.MAIN_QUEUE_TOPIC);
TenantProfileQueueConfiguration isolatedQueueConfiguration = createQueueConfig("IsolatedHighPriority", "tb_rule_engine.isolated_hp");
edgeTenantProfile.getProfileData().setQueueConfiguration(List.of(mainQueueConfiguration, isolatedQueueConfiguration));
// + 1 TenantProfile
// + 1 Queue main
// + 1 Queue isolated
edgeImitator.expectMessageAmount(3);
edgeTenantProfile = doPost("/api/tenantProfile", edgeTenantProfile, TenantProfile.class);
Assert.assertTrue(edgeImitator.waitForMessages());
Optional<TenantProfileUpdateMsg> tenantProfileUpdateMsgOpt = edgeImitator.findMessageByType(TenantProfileUpdateMsg.class);
Optional<TenantProfileUpdateMsg> tenantProfileUpdateMsgOpt = edgeImitator.findMessageByType(TenantProfileUpdateMsg.class);
Assert.assertTrue(tenantProfileUpdateMsgOpt.isPresent());
TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileUpdateMsgOpt.get();
TenantProfile tenantProfile = JacksonUtil.fromString(tenantProfileUpdateMsg.getEntity(), TenantProfile.class, true);
@ -96,7 +99,9 @@ public class TenantProfileEdgeTest extends AbstractEdgeTest {
loginTenantAdmin();
edgeImitator.expectMessageAmount(21);
// 25 sync message
// +1 isolated Queue
edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 1);
doPost("/api/edge/sync/" + edge.getId());
assertThat(edgeImitator.waitForMessages()).as("await for messages after edge sync rest api call").isTrue();

15
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -24,6 +24,7 @@ import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.junit.Assert;
import org.thingsboard.edge.rpc.EdgeGrpcClient;
import org.thingsboard.edge.rpc.EdgeRpcClient;
import org.thingsboard.server.controller.AbstractWebTest;
@ -386,7 +387,19 @@ public class EdgeImitator {
}
public boolean waitForMessages() throws InterruptedException {
return waitForMessages(AbstractWebTest.TIMEOUT);
boolean success = waitForMessages(AbstractWebTest.TIMEOUT);
if (!success) {
List<AbstractMessage> downlinkMsgs = getDownlinkMsgs();
for (AbstractMessage downlinkMsg : downlinkMsgs) {
log.error("{}\n{}", downlinkMsg.getClass(), downlinkMsg);
}
log.error("message count: {}", downlinkMsgs.size());
Assert.fail("Await for messages was not successful!");
}
return true;
}
public boolean waitForMessages(int timeoutInSeconds) throws InterruptedException {

10
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java

@ -44,7 +44,7 @@ import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.UUID;
@ -58,6 +58,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
@RequiredArgsConstructor
@Slf4j
public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao<EdgeEventEntity, EdgeEvent> implements EdgeEventDao {
private static final List<SortOrder> SORT_ORDERS = Collections.singletonList(new SortOrder("seqId"));
private final UUID systemTenantId = NULL_UUID;
@ -175,11 +176,6 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao<EdgeEventEnti
@Override
public PageData<EdgeEvent> findEdgeEvents(UUID tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) {
List<SortOrder> sortOrders = new ArrayList<>();
if (pageLink.getSortOrder() != null) {
sortOrders.add(pageLink.getSortOrder());
}
sortOrders.add(new SortOrder("seqId"));
return DaoUtil.toPageData(
edgeEventRepository
.findEdgeEventsByTenantIdAndEdgeId(
@ -190,7 +186,7 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao<EdgeEventEnti
pageLink.getEndTime(),
seqIdStart,
seqIdEnd,
DaoUtil.toPageable(pageLink, sortOrders)));
DaoUtil.toPageable(pageLink, SORT_ORDERS)));
}
@Override

25
dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java

@ -16,7 +16,6 @@
package org.thingsboard.server.dao.service;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.junit.Assert;
import org.junit.Before;
@ -38,8 +37,6 @@ import org.thingsboard.server.dao.edge.EdgeEventService;
import java.io.IOException;
import java.text.ParseException;
import java.util.ArrayList;
import java.util.List;
import static org.apache.commons.lang3.time.DateFormatUtils.ISO_8601_EXTENDED_DATETIME_FORMAT;
@ -103,26 +100,23 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
}
@Test
public void findEdgeEventsByTimeDescOrder() throws Exception {
public void findEdgeEventsBySeqIdOrder_createdTimeOrderIgnored() throws Exception {
EdgeId edgeId = new EdgeId(Uuids.timeBased());
DeviceId deviceId = new DeviceId(Uuids.timeBased());
List<ListenableFuture<Void>> futures = new ArrayList<>();
futures.add(saveEdgeEventWithProvidedTime(timeBeforeStartTime, edgeId, deviceId, tenantId));
futures.add(saveEdgeEventWithProvidedTime(eventTime, edgeId, deviceId, tenantId));
futures.add(saveEdgeEventWithProvidedTime(eventTime + 1, edgeId, deviceId, tenantId));
futures.add(saveEdgeEventWithProvidedTime(eventTime + 2, edgeId, deviceId, tenantId));
futures.add(saveEdgeEventWithProvidedTime(timeAfterEndTime, edgeId, deviceId, tenantId));
Futures.allAsList(futures).get();
saveEdgeEventWithProvidedTime(timeBeforeStartTime, edgeId, deviceId, tenantId).get();
saveEdgeEventWithProvidedTime(eventTime, edgeId, deviceId, tenantId).get();
saveEdgeEventWithProvidedTime(eventTime + 2, edgeId, deviceId, tenantId).get();
saveEdgeEventWithProvidedTime(eventTime + 1, edgeId, deviceId, tenantId).get();
saveEdgeEventWithProvidedTime(timeAfterEndTime, edgeId, deviceId, tenantId).get();
TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime);
PageData<EdgeEvent> edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, 0L, null, pageLink);
Assert.assertNotNull(edgeEvents.getData());
Assert.assertEquals(2, edgeEvents.getData().size());
Assert.assertEquals(Uuids.startOf(eventTime + 2), edgeEvents.getData().get(0).getUuidId());
Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().get(1).getUuidId());
Assert.assertEquals(Uuids.startOf(eventTime), edgeEvents.getData().get(0).getUuidId());
Assert.assertEquals(Uuids.startOf(eventTime + 2), edgeEvents.getData().get(1).getUuidId());
Assert.assertTrue(edgeEvents.hasNext());
Assert.assertNotNull(pageLink.nextPageLink());
@ -130,7 +124,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
Assert.assertNotNull(edgeEvents.getData());
Assert.assertEquals(1, edgeEvents.getData().size());
Assert.assertEquals(Uuids.startOf(eventTime), edgeEvents.getData().get(0).getUuidId());
Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().get(0).getUuidId());
Assert.assertFalse(edgeEvents.hasNext());
edgeEventDao.cleanupEvents(1);
@ -141,4 +135,5 @@ public class EdgeEventServiceTest extends AbstractServiceTest {
edgeEvent.setId(new EdgeEventId(Uuids.startOf(time)));
return edgeEventService.saveAsync(edgeEvent);
}
}
Loading…
Cancel
Save