Browse Source

Add check for negative and outdated disconnect event timestamp.

pull/9030/head
Dmytro Skarzhynets 3 years ago
committed by Dmytro Skarzhynets
parent
commit
653ea4a800
  1. 16
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  2. 90
      application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java

16
application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

@ -271,8 +271,19 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
if (cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId)) { if (cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId)) {
return; return;
} }
log.trace("[{}] on Device Disconnect [{}]", tenantId.getId(), deviceId.getId()); if (lastDisconnectTime < 0) {
log.trace("[{}][{}] On device disconnect: received negative last disconnect ts [{}]. Skipping this event.",
tenantId.getId(), deviceId.getId(), lastDisconnectTime);
return;
}
DeviceStateData stateData = getOrFetchDeviceStateData(deviceId); DeviceStateData stateData = getOrFetchDeviceStateData(deviceId);
long currentLastDisconnectTime = stateData.getState().getLastDisconnectTime();
if (lastDisconnectTime <= currentLastDisconnectTime) {
log.trace("[{}][{}] On device disconnect: received outdated last disconnect ts [{}]. Skipping this event. Current last disconnect ts [{}].",
tenantId.getId(), deviceId.getId(), lastDisconnectTime, currentLastDisconnectTime);
return;
}
log.trace("[{}][{}] On device disconnect: processing disconnect event with ts [{}].", tenantId.getId(), deviceId.getId(), lastDisconnectTime);
stateData.getState().setLastDisconnectTime(lastDisconnectTime); stateData.getState().setLastDisconnectTime(lastDisconnectTime);
save(deviceId, LAST_DISCONNECT_TIME, lastDisconnectTime); save(deviceId, LAST_DISCONNECT_TIME, lastDisconnectTime);
pushRuleEngineMessage(stateData, TbMsgType.DISCONNECT_EVENT); pushRuleEngineMessage(stateData, TbMsgType.DISCONNECT_EVENT);
@ -302,8 +313,6 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
tenantId.getId(), deviceId.getId(), lastInactivityTime); tenantId.getId(), deviceId.getId(), lastInactivityTime);
return; return;
} }
log.trace("[{}][{}] On device inactivity: processing inactivity event with ts [{}].",
tenantId.getId(), deviceId.getId(), lastInactivityTime);
DeviceStateData stateData = getOrFetchDeviceStateData(deviceId); DeviceStateData stateData = getOrFetchDeviceStateData(deviceId);
long currentLastInactivityAlarmTime = stateData.getState().getLastInactivityAlarmTime(); long currentLastInactivityAlarmTime = stateData.getState().getLastInactivityAlarmTime();
if (lastInactivityTime <= currentLastInactivityAlarmTime) { if (lastInactivityTime <= currentLastInactivityAlarmTime) {
@ -311,6 +320,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
tenantId.getId(), deviceId.getId(), lastInactivityTime, currentLastInactivityAlarmTime); tenantId.getId(), deviceId.getId(), lastInactivityTime, currentLastInactivityAlarmTime);
return; return;
} }
log.trace("[{}][{}] On device inactivity: processing inactivity event with ts [{}].", tenantId.getId(), deviceId.getId(), lastInactivityTime);
reportInactivity(lastInactivityTime, deviceId, stateData); reportInactivity(lastInactivityTime, deviceId, stateData);
} }

90
application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java

@ -83,6 +83,7 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACT
import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_ALARM_TIME; import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_ALARM_TIME;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_TIMEOUT; import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_TIMEOUT;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_ACTIVITY_TIME; import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_ACTIVITY_TIME;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_DISCONNECT_TIME;
@ExtendWith(MockitoExtension.class) @ExtendWith(MockitoExtension.class)
public class DefaultDeviceStateServiceTest { public class DefaultDeviceStateServiceTest {
@ -125,6 +126,91 @@ public class DefaultDeviceStateServiceTest {
tpi = TopicPartitionInfo.builder().myPartition(true).build(); tpi = TopicPartitionInfo.builder().myPartition(true).build();
} }
@Test
public void givenDeviceBelongsToExternalPartition_whenOnDeviceDisconnect_thenCleansStateAndDoesNotReportDisconnect() {
// GIVEN
doReturn(true).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
// WHEN
service.onDeviceDisconnect(tenantId, deviceId, System.currentTimeMillis());
// THEN
then(service).should().cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
then(service).should(never()).getOrFetchDeviceStateData(deviceId);
then(clusterService).shouldHaveNoInteractions();
then(notificationRuleProcessor).shouldHaveNoInteractions();
then(telemetrySubscriptionService).shouldHaveNoInteractions();
}
@ParameterizedTest
@ValueSource(longs = {Long.MIN_VALUE, -100, -1})
public void givenNegativeLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long negativeLastDisconnectTime) {
// GIVEN
doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
// WHEN
service.onDeviceDisconnect(tenantId, deviceId, negativeLastDisconnectTime);
// THEN
then(service).should(never()).getOrFetchDeviceStateData(deviceId);
then(clusterService).shouldHaveNoInteractions();
then(notificationRuleProcessor).shouldHaveNoInteractions();
then(telemetrySubscriptionService).shouldHaveNoInteractions();
}
@ParameterizedTest
@MethodSource("provideOutdatedTimestamps")
public void givenOutdatedLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long outdatedLastDisconnectTime, long currentLastDisconnectTime) {
// GIVEN
doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
var deviceStateData = DeviceStateData.builder()
.tenantId(tenantId)
.deviceId(deviceId)
.state(DeviceState.builder().lastDisconnectTime(currentLastDisconnectTime).build())
.build();
service.deviceStates.put(deviceId, deviceStateData);
// WHEN
service.onDeviceDisconnect(tenantId, deviceId, outdatedLastDisconnectTime);
// THEN
then(clusterService).shouldHaveNoInteractions();
then(notificationRuleProcessor).shouldHaveNoInteractions();
then(telemetrySubscriptionService).shouldHaveNoInteractions();
}
@Test
public void givenDeviceBelongsToMyPartition_whenOnDeviceDisconnect_thenReportsDisconnect() {
// GIVEN
var deviceStateData = DeviceStateData.builder()
.tenantId(tenantId)
.deviceId(deviceId)
.state(DeviceState.builder().build())
.metaData(new TbMsgMetaData())
.build();
doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
service.deviceStates.put(deviceId, deviceStateData);
long lastDisconnectTime = System.currentTimeMillis();
// WHEN
service.onDeviceDisconnect(tenantId, deviceId, lastDisconnectTime);
// THEN
then(telemetrySubscriptionService).should().saveAttrAndNotify(
eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(DataConstants.SERVER_SCOPE),
eq(LAST_DISCONNECT_TIME), eq(lastDisconnectTime), any()
);
var msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
then(clusterService).should().pushMsgToRuleEngine(eq(tenantId), eq(deviceId), msgCaptor.capture(), any());
var actualMsg = msgCaptor.getValue();
assertThat(actualMsg.getType()).isEqualTo(TbMsgType.DISCONNECT_EVENT.name());
assertThat(actualMsg.getOriginator()).isEqualTo(deviceId);
}
@Test @Test
public void givenDeviceBelongsToExternalPartition_whenOnDeviceInactivity_thenCleansStateAndDoesNotReportInactivity() { public void givenDeviceBelongsToExternalPartition_whenOnDeviceInactivity_thenCleansStateAndDoesNotReportInactivity() {
// GIVEN // GIVEN
@ -158,7 +244,7 @@ public class DefaultDeviceStateServiceTest {
} }
@ParameterizedTest @ParameterizedTest
@MethodSource @MethodSource("provideOutdatedTimestamps")
public void givenOutdatedLastInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent(long outdatedLastActivityTime, long currentLastActivityTime) { public void givenOutdatedLastInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent(long outdatedLastActivityTime, long currentLastActivityTime) {
// GIVEN // GIVEN
doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId);
@ -179,7 +265,7 @@ public class DefaultDeviceStateServiceTest {
then(telemetrySubscriptionService).shouldHaveNoInteractions(); then(telemetrySubscriptionService).shouldHaveNoInteractions();
} }
private static Stream<Arguments> givenOutdatedLastInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent() { private static Stream<Arguments> provideOutdatedTimestamps() {
return Stream.of( return Stream.of(
Arguments.of(0, 0), Arguments.of(0, 0),
Arguments.of(0, 100), Arguments.of(0, 100),

Loading…
Cancel
Save