Browse Source

Merge pull request #15305 from dashevchenko/activityEventOnDeviceDisconnectOnExpiration

No activity reporting on device disconnect event
pull/14719/merge
Viacheslav Klimov 4 months ago
committed by GitHub
parent
commit
0d29f60a62
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  2. 58
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java

4
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -508,7 +508,9 @@ public class DefaultTransportService extends TransportActivityManager implements
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
recordActivityInternal(sessionInfo);
if (msg.getEvent() != TransportProtos.SessionEvent.CLOSED) {
recordActivityInternal(sessionInfo);
}
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setSessionEvent(msg).build(), callback);
}

58
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java

@ -26,12 +26,14 @@ import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.activity.ActivityReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategyType;
import org.thingsboard.server.common.transport.limits.TransportRateLimitService;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.UUID;
@ -41,7 +43,12 @@ import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.ArgumentMatchers.nullable;
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@ -514,4 +521,55 @@ public class TransportActivityManagerTest {
verify(transportServiceMock, never()).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
}
@Test
void givenSessionClosedEvent_whenProcessingSessionEvent_thenShouldNotRecordActivity() {
// GIVEN - simulates the session-expiry path: session already removed from sessions map,
// then process(SESSION_CLOSED) is called. Activity must NOT be recorded to avoid active
// status update from false to true
var rateLimitServiceMock = mock(TransportRateLimitService.class);
ReflectionTestUtils.setField(transportServiceMock, "rateLimitService", rateLimitServiceMock);
when(rateLimitServiceMock.checkLimits(any(), nullable(DeviceId.class), any(), anyInt(), anyBoolean())).thenReturn(null);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
doCallRealMethod().when(transportServiceMock).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
// WHEN
transportServiceMock.process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
// THEN
verify(transportServiceMock, never()).onActivity(any(), any(), anyLong());
verify(transportServiceMock).sendToDeviceActor(eq(sessionInfo), any(), isNull());
}
@Test
void givenSessionOpenEvent_whenProcessingSessionEvent_thenShouldRecordActivity() {
// GIVEN
var rateLimitServiceMock = mock(TransportRateLimitService.class);
ReflectionTestUtils.setField(transportServiceMock, "rateLimitService", rateLimitServiceMock);
when(rateLimitServiceMock.checkLimits(any(), nullable(DeviceId.class), any(), anyInt(), anyBoolean())).thenReturn(null);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
var sessionOpenMsg = TransportProtos.SessionEventMsg.newBuilder()
.setSessionType(TransportProtos.SessionType.ASYNC)
.setEvent(TransportProtos.SessionEvent.OPEN)
.build();
when(transportServiceMock.toSessionId(sessionInfo)).thenReturn(SESSION_ID);
long expectedActivityTime = 100L;
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(expectedActivityTime);
doCallRealMethod().when(transportServiceMock).process(sessionInfo, sessionOpenMsg, null);
// WHEN
transportServiceMock.process(sessionInfo, sessionOpenMsg, null);
// THEN
verify(transportServiceMock).onActivity(SESSION_ID, sessionInfo, expectedActivityTime);
verify(transportServiceMock).sendToDeviceActor(eq(sessionInfo), any(), isNull());
}
}

Loading…
Cancel
Save