diff --git a/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java b/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java index 6748bced4d..fdb1710189 100644 --- a/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java +++ b/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java @@ -81,7 +81,7 @@ public class IntegrationActivityManagerTest { private DefaultPlatformIntegrationService integrationServiceMock; @Test - void testReportActivity() { + void givenKeyAndTimeToReport_whenReportingActivity_thenShouldCorrectlyReportActivity() { // GIVEN var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); @@ -184,58 +184,54 @@ public class IntegrationActivityManagerTest { @ParameterizedTest @MethodSource("provideTestParamsForCreateNewState") void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy( - String reportingStrategyName, Class reportingStrategyClass + String reportingStrategyName, ActivityStrategy reportingStrategy ) { // GIVEN var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); when(integrationServiceMock.createNewState(key)).thenCallRealMethod(); ReflectionTestUtils.setField(integrationServiceMock, "reportingStrategyName", reportingStrategyName); + ActivityState expectedState = new ActivityState<>(); + expectedState.setStrategy(reportingStrategy); + // WHEN - ActivityState newState = integrationServiceMock.createNewState(key); + ActivityState actualState = integrationServiceMock.createNewState(key); // THEN - assertThat(newState).isNotNull(); - assertThat(newState.getLastRecordedTime()).isEqualTo(0L); - assertThat(newState.getLastReportedTime()).isEqualTo(0L); - assertThat(newState.getStrategy()).isInstanceOf(reportingStrategyClass); + assertThat(actualState).isEqualTo(expectedState); } private static Stream provideTestParamsForCreateNewState() { return Stream.of( - Arguments.of("ALL", AllEventsActivityStrategy.class), - Arguments.of("FIRST", FirstEventActivityStrategy.class), - Arguments.of("LAST", LastEventActivityStrategy.class), - Arguments.of("FIRST_AND_LAST", FirstAndLastEventActivityStrategy.class) + Arguments.of("ALL", new AllEventsActivityStrategy()), + Arguments.of("FIRST", new FirstEventActivityStrategy()), + Arguments.of("LAST", new LastEventActivityStrategy()), + Arguments.of("FIRST_AND_LAST", new FirstAndLastEventActivityStrategy()) ); } @Test - void givenActivityState_whenUpdatingActivityState_thenShouldReturnSameInstanceWithNoChanges() { + void givenActivityState_whenUpdatingActivityState_thenShouldReturnSameInstanceWithNoInteractions() { // GIVEN var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); - long expectedLastRecordedTime = 123L; - long expectedLastReportedTime = 312L; - ActivityStrategy expectedStrategy = spy(new FirstEventActivityStrategy()); + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); - ActivityState expectedState = new ActivityState<>(); - expectedState.setLastRecordedTime(expectedLastRecordedTime); - expectedState.setLastReportedTime(expectedLastReportedTime); - expectedState.setStrategy(expectedStrategy); + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(123L); + state.setLastReportedTime(312L); + state.setStrategy(strategySpy); + ActivityState stateSpy = spy(state); - when(integrationServiceMock.updateState(key, expectedState)).thenCallRealMethod(); + when(integrationServiceMock.updateState(key, state)).thenCallRealMethod(); // WHEN - ActivityState actualNewState = integrationServiceMock.updateState(key, expectedState); + ActivityState updatedState = integrationServiceMock.updateState(key, state); // THEN - assertThat(actualNewState).isSameAs(expectedState); - assertThat(actualNewState.getLastRecordedTime()).isEqualTo(expectedLastRecordedTime); - assertThat(actualNewState.getLastReportedTime()).isEqualTo(expectedLastReportedTime); - assertThat(actualNewState.getStrategy()).isSameAs(expectedStrategy); - verifyNoInteractions(expectedStrategy); - assertThat(actualNewState.getMetadata()).isNull(); + assertThat(updatedState).isSameAs(state); + verifyNoInteractions(stateSpy); + verifyNoInteractions(strategySpy); } @ParameterizedTest @@ -258,7 +254,8 @@ public class IntegrationActivityManagerTest { return Stream.of( Arguments.of(10L, 0L, 9L), Arguments.of(10L, 7L, 2L), - Arguments.of(10L, 8L, 1L) + Arguments.of(10L, 8L, 1L), + Arguments.of(10000L, 5000L, 3000L) ); } @@ -282,12 +279,13 @@ public class IntegrationActivityManagerTest { return Stream.of( Arguments.of(10L, 9L, 2L), Arguments.of(10L, 0L, 11L), - Arguments.of(10L, 8L, 3L) + Arguments.of(10L, 8L, 3L), + Arguments.of(10000L, 8000L, 3000L) ); } @Test - void givenKeyAndVoidMetadata_whenOnStateExpiryCalled_thenShouldDoNothing() { + void givenKeyAndMetadata_whenOnStateExpiryCalled_thenShouldDoNothing() { // GIVEN var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); doCallRealMethod().when(integrationServiceMock).onStateExpiry(key, null); diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index 1e68779b70..735ec1d524 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -212,7 +212,7 @@ public class DefaultCoapClientContext implements CoapClientContext { public void reportActivity() { for (TbCoapClientState state : clients.values()) { if (state.getSession() != null) { - transportService.reportActivity(state.getSession()); + transportService.recordActivity(state.getSession()); } } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 90bddd2edc..8360111a54 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -322,7 +322,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement case PINGREQ: if (checkConnected(ctx, msg)) { ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + transportService.recordActivity(deviceSessionCtx.getSessionInfo()); } break; case DISCONNECT: @@ -350,7 +350,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) { if (gatewaySessionHandler != null) { handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg); - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + transportService.recordActivity(deviceSessionCtx.getSessionInfo()); } else { log.error("[gatewaySessionHandler] is null, [{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId); } @@ -526,7 +526,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); attrReqTopicType = TopicType.V2; } else { - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + transportService.recordActivity(deviceSessionCtx.getSessionInfo()); ack(ctx, msgId, ReturnCode.TOPIC_NAME_INVALID); } } catch (AdaptorException e) { @@ -796,7 +796,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } if (!activityReported) { - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + transportService.recordActivity(deviceSessionCtx.getSessionInfo()); } ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); } @@ -894,7 +894,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } if (!activityReported) { - transportService.reportActivity(deviceSessionCtx.getSessionInfo()); + transportService.recordActivity(deviceSessionCtx.getSessionInfo()); } ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId(), unSubResults)); } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java index 5d41918cab..3a8811df97 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java @@ -416,7 +416,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde } private void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { - transportService.reportActivity(sessionInfo); + transportService.recordActivity(sessionInfo); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index f51b399a5b..d394af3101 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -145,7 +145,7 @@ public interface TransportService { SessionMetaData registerSyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout); - void reportActivity(SessionInfoProto sessionInfo); + void recordActivity(SessionInfoProto sessionInfo); void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java index b5f4bf8f9e..07217e57e8 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java @@ -129,7 +129,8 @@ public abstract class AbstractActivityManager implements Activity } } - protected long getLastRecordedTime(Key key) { + @Override + public long getLastRecordedTime(Key key) { ActivityState state = states.get(key); return state == null ? 0L : state.getLastRecordedTime(); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java index f19b9e30ca..847e8f4472 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java @@ -36,4 +36,6 @@ public interface ActivityManager { void onActivity(Key key); + long getLastRecordedTime(Key key); + } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java index 27797c5315..ba7ceab517 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java @@ -30,6 +30,9 @@ */ package org.thingsboard.server.common.transport.activity.strategy; +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode public class AllEventsActivityStrategy implements ActivityStrategy { @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java index c2605cdeaf..225d02ccb8 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java @@ -30,6 +30,9 @@ */ package org.thingsboard.server.common.transport.activity.strategy; +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode public class FirstAndLastEventActivityStrategy implements ActivityStrategy { private boolean firstEventReceived; diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java index d61c7a68aa..f4f53de6c0 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java @@ -30,6 +30,9 @@ */ package org.thingsboard.server.common.transport.activity.strategy; +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode public class FirstEventActivityStrategy implements ActivityStrategy { private boolean firstEventReceived; diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java index e22263096e..f16fbed8d6 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java @@ -30,6 +30,9 @@ */ package org.thingsboard.server.common.transport.activity.strategy; +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode public class LastEventActivityStrategy implements ActivityStrategy { @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 9cd41e7a62..da0ef3ca99 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -558,7 +558,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) .setSessionEvent(msg).build(), callback); } @@ -578,7 +578,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); TenantId tenantId = getTenantId(sessionInfo); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); @@ -639,7 +639,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) .setGetAttributes(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -652,7 +652,7 @@ public class DefaultTransportService extends AbstractActivityManager(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -665,7 +665,7 @@ public class DefaultTransportService extends AbstractActivityManager(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -674,7 +674,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setToDeviceRPCCallResponse(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -683,7 +683,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback); } } @@ -704,7 +704,7 @@ public class DefaultTransportService extends AbstractActivityManager(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY)); @@ -736,7 +736,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); UUID sessionId = toSessionId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo); DeviceId deviceId = getDeviceId(sessionInfo); @@ -761,7 +761,7 @@ public class DefaultTransportService extends AbstractActivityManager callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) .setClaimDevice(msg).build(), callback); } @@ -780,11 +780,11 @@ public class DefaultTransportService extends AbstractActivityManager lastRecordedTime; + return (getCurrentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime; + } + + long getCurrentTimeMillis() { + return System.currentTimeMillis(); } @Override @@ -937,7 +941,7 @@ public class DefaultTransportService extends AbstractActivityManager sessions; + + @BeforeEach + public void setup() { + sessions = new ConcurrentHashMap<>(); + ReflectionTestUtils.setField(transportServiceMock, "sessions", sessions); + } + +// @Override +// protected void reportActivity(UUID sessionId, TransportProtos.SessionInfoProto currentSessionInfo, long timeToReport, ActivityReportCallback callback) { +// log.debug("[{}] Reporting activity state for session with id: [{}]. Time to report: [{}].", name, sessionId, timeToReport); +// SessionMetaData session = sessions.get(sessionId); +// TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() +// .setAttributeSubscription(session != null && session.isSubscribedToAttributes()) +// .setRpcSubscription(session != null && session.isSubscribedToRPC()) +// .setLastActivityTime(timeToReport) +// .build(); +// TransportProtos.SessionInfoProto sessionInfo = session != null ? session.getSessionInfo() : currentSessionInfo; +// process(sessionInfo, subscriptionInfo, new TransportServiceCallback<>() { +// @Override +// public void onSuccess(Void msgAcknowledged) { +// callback.onSuccess(sessionId, timeToReport); +// +// } +// +// @Override +// public void onError(Throwable e) { +// callback.onFailure(sessionId, e); +// } +// }); +// } + + @Test + void givenKeyAndTimeToReportAndSessionExists_whenReportingActivity_thenShouldReportActivityWithSubscriptionsAndSessionInfoFromSession() { + // GIVEN + long expectedTime = 123L; + boolean expectedAttributesSubscription = true; + boolean expectedRPCSubscription = true; + TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.getDefaultInstance(); + + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + SessionMetaData session = new SessionMetaData(expectedSessionInfo, TransportProtos.SessionType.ASYNC, listenerMock); + session.setSubscribedToAttributes(expectedAttributesSubscription); + session.setSubscribedToRPC(expectedRPCSubscription); + sessions.put(SESSION_ID, session); + + ActivityReportCallback callbackMock = mock(ActivityReportCallback.class); + + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + + doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock); + + // WHEN + transportServiceMock.reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock); + + // THEN + ArgumentCaptor sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class); + ArgumentCaptor subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class); + ArgumentCaptor> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class); + + verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture()); + + assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo); + + TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() + .setAttributeSubscription(expectedAttributesSubscription) + .setRpcSubscription(expectedRPCSubscription) + .setLastActivityTime(expectedTime) + .build(); + assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo); + + TransportServiceCallback queueCallback = callbackCaptor.getValue(); + + queueCallback.onSuccess(null); + verify(callbackMock).onSuccess(SESSION_ID, expectedTime); + + var throwable = new Throwable(); + queueCallback.onError(throwable); + verify(callbackMock).onFailure(SESSION_ID, throwable); + } + + @Test + void givenKeyAndTimeToReportAndSessionDoesNotExist_whenReportingActivity_thenShouldReportActivityWithNoSubscriptionsAndPreviousSessionInfo() { + // GIVEN + long expectedTime = 123L; + boolean expectedAttributesSubscription = false; + boolean expectedRPCSubscription = false; + TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + + ActivityReportCallback callbackMock = mock(ActivityReportCallback.class); + + doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock); + + // WHEN + transportServiceMock.reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock); + + // THEN + ArgumentCaptor sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class); + ArgumentCaptor subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class); + ArgumentCaptor> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class); + + verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture()); + + assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo); + + TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() + .setAttributeSubscription(expectedAttributesSubscription) + .setRpcSubscription(expectedRPCSubscription) + .setLastActivityTime(expectedTime) + .build(); + assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo); + + TransportServiceCallback queueCallback = callbackCaptor.getValue(); + + queueCallback.onSuccess(null); + verify(callbackMock).onSuccess(SESSION_ID, expectedTime); + + var throwable = new Throwable(); + queueCallback.onError(throwable); + verify(callbackMock).onFailure(SESSION_ID, throwable); + } + + @Test + void givenActivityHappened_whenRecordActivity_thenShouldDelegateToOnActivity() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + doCallRealMethod().when(transportServiceMock).recordActivity(sessionInfo); + when(transportServiceMock.toSessionId(sessionInfo)).thenReturn(SESSION_ID); + + // WHEN + transportServiceMock.recordActivity(sessionInfo); + + // THEN + verify(transportServiceMock).onActivity(SESSION_ID); + } + + @ParameterizedTest + @MethodSource("provideTestParamsForCreateNewState") + void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy( + String reportingStrategyName, ActivityStrategy reportingStrategy + ) { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + ReflectionTestUtils.setField(transportServiceMock, "reportingStrategyName", reportingStrategyName); + + when(transportServiceMock.createNewState(SESSION_ID)).thenCallRealMethod(); + + ActivityState expectedState = new ActivityState<>(); + expectedState.setStrategy(reportingStrategy); + expectedState.setMetadata(sessionInfo); + + // WHEN + ActivityState actualState = transportServiceMock.createNewState(SESSION_ID); + + // THEN + assertThat(actualState).isEqualTo(expectedState); + } + + private static Stream provideTestParamsForCreateNewState() { + return Stream.of( + Arguments.of("ALL", new AllEventsActivityStrategy()), + Arguments.of("FIRST", new FirstEventActivityStrategy()), + Arguments.of("LAST", new LastEventActivityStrategy()), + Arguments.of("FIRST_AND_LAST", new FirstAndLastEventActivityStrategy()) + ); + } + + @Test + void givenSessionDoesNotExist_whenUpdatingActivityState_thenShouldReturnNull() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(123L); + state.setLastReportedTime(312L); + state.setMetadata(sessionInfo); + state.setStrategy(new FirstEventActivityStrategy()); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isNull(); + } + + @Test + void givenNoGwSessionId_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + long lastRecordedTime = 123L; + long lastReportedTime = 312L; + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setLastReportedTime(lastReportedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + state.setStrategy(strategySpy); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isSameAs(state); + assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); + assertThat(updatedState.getLastReportedTime()).isEqualTo(lastReportedTime); + assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); + assertThat(updatedState.getStrategy()).isEqualTo(strategySpy); + verifyNoInteractions(strategySpy); + } + + @Test + void givenHasGwSessionIdButGwSessionIsNotNull_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { + // GIVEN + var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + long lastRecordedTime = 123L; + long lastReportedTime = 312L; + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setLastReportedTime(lastReportedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + state.setStrategy(strategySpy); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isSameAs(state); + assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); + assertThat(updatedState.getLastReportedTime()).isEqualTo(lastReportedTime); + assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); + assertThat(updatedState.getStrategy()).isEqualTo(strategySpy); + verifyNoInteractions(strategySpy); + + verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId); + } + + @Test + void givenHasGwSessionWithoutOverwriteEnabled_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { + // GIVEN + var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); + TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); + sessions.put(gwSessionId, new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock)); + + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + long lastRecordedTime = 123L; + long lastReportedTime = 312L; + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setLastReportedTime(lastReportedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + state.setStrategy(strategySpy); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isSameAs(state); + assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); + assertThat(updatedState.getLastReportedTime()).isEqualTo(lastReportedTime); + assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); + assertThat(updatedState.getStrategy()).isEqualTo(strategySpy); + verifyNoInteractions(strategySpy); + + verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId); + } + + @Test + void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsGreater_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoAndLastRecordedTime() { + // GIVEN + var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); + TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); + SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock); + gwSession.setOverwriteActivityTime(true); + sessions.put(gwSessionId, gwSession); + + long gwLastRecordedTime = 500L; + when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime); + + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + long lastRecordedTime = 123L; + long lastReportedTime = 312L; + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setLastReportedTime(lastReportedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + state.setStrategy(strategySpy); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isSameAs(state); + assertThat(updatedState.getLastRecordedTime()).isEqualTo(gwLastRecordedTime); + assertThat(updatedState.getLastReportedTime()).isEqualTo(lastReportedTime); + assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); + assertThat(updatedState.getStrategy()).isEqualTo(strategySpy); + verifyNoInteractions(strategySpy); + } + + @Test + void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsLess_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoOnly() { + // GIVEN + var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); + TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); + SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock); + gwSession.setOverwriteActivityTime(true); + sessions.put(gwSessionId, gwSession); + + long gwLastRecordedTime = 100L; + when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime); + + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) + .setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + + long lastRecordedTime = 123L; + long lastReportedTime = 312L; + ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy()); + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setLastReportedTime(lastReportedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + state.setStrategy(strategySpy); + + when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); + + // WHEN + ActivityState updatedState = transportServiceMock.updateState(SESSION_ID, state); + + // THEN + assertThat(updatedState).isSameAs(state); + assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); + assertThat(updatedState.getLastReportedTime()).isEqualTo(lastReportedTime); + assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); + assertThat(updatedState.getStrategy()).isEqualTo(strategySpy); + verifyNoInteractions(strategySpy); + } + + @ParameterizedTest + @MethodSource("provideTestParamsForHasExpiredTrue") + public void givenExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnTrue(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) { + // GIVEN + ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout); + + when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis); + when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod(); + + // WHEN + boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime); + + // THEN + assertThat(hasExpired).isTrue(); + } + + private static Stream provideTestParamsForHasExpiredTrue() { + return Stream.of( + Arguments.of(10L, 0L, 9L), + Arguments.of(10L, 7L, 2L), + Arguments.of(10L, 8L, 1L), + Arguments.of(10000L, 5000L, 3000L) + ); + } + + @ParameterizedTest + @MethodSource("provideTestParamsForHasExpiredFalse") + public void givenNotExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnFalse(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) { + // GIVEN + ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout); + + when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis); + when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod(); + + // WHEN + boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime); + + // THEN + assertThat(hasExpired).isFalse(); + } + + private static Stream provideTestParamsForHasExpiredFalse() { + return Stream.of( + Arguments.of(10L, 9L, 2L), + Arguments.of(10L, 0L, 11L), + Arguments.of(10L, 8L, 3L), + Arguments.of(10000L, 8000L, 3000L) + ); + } + + @Test + void givenSessionExists_whenOnStateExpiryCalled_thenShouldPerformExpirationActions() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + SessionMsgListener listenerMock = mock(SessionMsgListener.class); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); + doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo); + + // WHEN + transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo); + + // THEN + assertThat(sessions.containsKey(SESSION_ID)).isFalse(); + verify(transportServiceMock).deregisterSession(sessionInfo); + verify(transportServiceMock).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); + verify(listenerMock).onRemoteSessionCloseCommand(SESSION_ID, SESSION_EXPIRED_NOTIFICATION_PROTO); + } + + @Test + void givenSessionDoesNotExist_whenOnStateExpiryCalled_thenShouldNotPerformExpirationActions() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo); + + // WHEN + transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo); + + // THEN + assertThat(sessions.containsKey(SESSION_ID)).isFalse(); + verify(transportServiceMock, never()).deregisterSession(sessionInfo); + verify(transportServiceMock, never()).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); + } + +}