diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 37070718a3..e010cf3050 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -185,6 +185,10 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceStates = new ConcurrentHashMap<>(); @@ -803,7 +807,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, key, value)); + telemetryTtl, new TelemetrySaveCallback<>(deviceId, key, value)); } else { tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); } @@ -814,7 +818,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, key, value)); + telemetryTtl, new TelemetrySaveCallback<>(deviceId, key, value)); } else { tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index e746dc9dba..2f710f3e92 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -785,6 +785,10 @@ state: # If 'persistToTelemetry' is changed from 'false' to 'true': 'CREATE OR REPLACE VIEW device_info_view AS SELECT * FROM device_info_active_ts_view;' # If 'persistToTelemetry' is changed from 'true' to 'false': 'CREATE OR REPLACE VIEW device_info_view AS SELECT * FROM device_info_active_attribute_view;' persistToTelemetry: "${PERSIST_STATE_TO_TELEMETRY:false}" + # Millisecond value defining time-to-live for device state telemetry data (e.g. 'active', 'lastActivityTime'). + # Used only when state.persistToTelemetry is set to 'true' and Cassandra is used for timeseries data. + # 0 means time-to-live mechanism is disabled. + telemetryTtl: "${STATE_TELEMETRY_TTL:0}" # Tbel parameters tbel: @@ -866,6 +870,14 @@ transport: inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" # Interval of periodic check for expired sessions and report of the changes to session last activity time report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:3000}" + activity: + # This property specifies the strategy for reporting activity events within each reporting period. + # The accepted values are 'FIRST', 'LAST', 'FIRST_AND_LAST' and 'ALL'. + # - 'FIRST': Only the first activity event in each reporting period is reported. + # - 'LAST': Only the last activity event in the reporting period is reported. + # - 'FIRST_AND_LAST': Both the first and last activity events in the reporting period are reported. + # - 'ALL': All activity events in the reporting period are reported. + reporting_strategy: "${TB_TRANSPORT_ACTIVITY_REPORTING_STRATEGY:LAST}" json: # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" @@ -1695,4 +1707,4 @@ management: web: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). - include: '${METRICS_ENDPOINTS_EXPOSE:info}' \ No newline at end of file + include: '${METRICS_ENDPOINTS_EXPOSE:info}' 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 new file mode 100644 index 0000000000..84d6d37824 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java @@ -0,0 +1,187 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity; + +import lombok.Data; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; +import org.thingsboard.server.queue.scheduler.SchedulerComponent; + +import java.util.Map; +import java.util.Random; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +@Slf4j +public abstract class AbstractActivityManager implements ActivityManager { + + private final ConcurrentMap states = new ConcurrentHashMap<>(); + + @Autowired + protected SchedulerComponent scheduler; + + @Data + private class ActivityStateWrapper { + + private volatile ActivityState state; + private volatile long lastReportedTime; + private volatile ActivityStrategy strategy; + + } + + protected void init() { + var reportingPeriodMillis = getReportingPeriodMillis(); + scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodMillis), reportingPeriodMillis, TimeUnit.MILLISECONDS); + } + + protected abstract long getReportingPeriodMillis(); + + protected abstract ActivityState createNewState(Key key); + + protected abstract ActivityStrategy getStrategy(); + + protected abstract ActivityState updateState(Key key, ActivityState state); + + protected abstract boolean hasExpired(long lastRecordedTime); + + protected abstract void onStateExpiry(Key key, Metadata metadata); + + protected abstract void reportActivity(Key key, Metadata metadata, long timeToReport, ActivityReportCallback callback); + + @Override + public void onActivity(Key key, long newLastRecordedTime) { + if (key == null) { + log.error("Failed to process activity event: provided activity key is null."); + return; + } + log.debug("Received activity event for key: [{}]", key); + + var shouldReport = new AtomicBoolean(false); + var lastRecordedTime = new AtomicLong(); + var lastReportedTime = new AtomicLong(); + var metadata = new AtomicReference(); + + var activityStateWrapper = states.compute(key, (__, stateWrapper) -> { + if (stateWrapper == null) { + var newState = createNewState(key); + if (newState == null) { + return null; + } + stateWrapper = new ActivityStateWrapper(); + stateWrapper.setState(newState); + stateWrapper.setStrategy(getStrategy()); + } + var state = stateWrapper.getState(); + if (state.getLastRecordedTime() < newLastRecordedTime) { + state.setLastRecordedTime(newLastRecordedTime); + } + shouldReport.set(stateWrapper.getStrategy().onActivity()); + lastRecordedTime.set(state.getLastRecordedTime()); + lastReportedTime.set(stateWrapper.getLastReportedTime()); + metadata.set(state.getMetadata()); + return stateWrapper; + }); + + if (activityStateWrapper == null) { + return; + } + + if (shouldReport.get() && lastReportedTime.get() < lastRecordedTime.get()) { + log.debug("Going to report first activity event for key: [{}].", key); + reportActivity(key, metadata.get(), lastRecordedTime.get(), new ActivityReportCallback<>() { + @Override + public void onSuccess(Key key, long reportedTime) { + updateLastReportedTime(key, reportedTime); + } + + @Override + public void onFailure(Key key, Throwable t) { + log.debug("Failed to report first activity event for key: [{}].", key, t); + } + }); + } + } + + @Override + public void onReportingPeriodEnd() { + log.debug("Going to end reporting period."); + for (Map.Entry entry : states.entrySet()) { + var key = entry.getKey(); + var stateWrapper = entry.getValue(); + var currentState = stateWrapper.getState(); + + long lastRecordedTime = currentState.getLastRecordedTime(); + long lastReportedTime = stateWrapper.getLastReportedTime(); + var metadata = currentState.getMetadata(); + + boolean hasExpired; + boolean shouldReport; + + var updatedState = updateState(key, currentState); + if (updatedState != null) { + stateWrapper.setState(updatedState); + lastRecordedTime = updatedState.getLastRecordedTime(); + metadata = updatedState.getMetadata(); + hasExpired = hasExpired(lastRecordedTime); + shouldReport = stateWrapper.getStrategy().onReportingPeriodEnd(); + } else { + states.remove(key); + hasExpired = false; + shouldReport = true; + } + + if (hasExpired) { + states.remove(key); + onStateExpiry(key, metadata); + shouldReport = true; + } + + if (shouldReport && lastReportedTime < lastRecordedTime) { + log.debug("Going to report last activity event for key: [{}].", key); + reportActivity(key, metadata, lastRecordedTime, new ActivityReportCallback<>() { + @Override + public void onSuccess(Key key, long reportedTime) { + updateLastReportedTime(key, reportedTime); + } + + @Override + public void onFailure(Key key, Throwable t) { + log.debug("Failed to report last activity event for key: [{}].", key, t); + } + }); + } + } + } + + @Override + public long getLastRecordedTime(Key key) { + ActivityStateWrapper stateWrapper = states.get(key); + return stateWrapper == null ? 0L : stateWrapper.getState().getLastRecordedTime(); + } + + private void updateLastReportedTime(Key key, long newLastReportedTime) { + states.computeIfPresent(key, (__, stateWrapper) -> { + stateWrapper.setLastReportedTime(Math.max(stateWrapper.getLastReportedTime(), newLastReportedTime)); + return stateWrapper; + }); + } + +} 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 new file mode 100644 index 0000000000..0f2145ab6f --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity; + +public interface ActivityManager { + + void onActivity(Key key, long activityTimeMillis); + + void onReportingPeriodEnd(); + + long getLastRecordedTime(Key key); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityReportCallback.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityReportCallback.java new file mode 100644 index 0000000000..c58f2bfa34 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityReportCallback.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity; + +public interface ActivityReportCallback { + + void onSuccess(Key key, long reportedTime); + + void onFailure(Key key, Throwable t); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionActivityData.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java similarity index 52% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionActivityData.java rename to common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java index 3e0be50ce6..bde5e05645 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionActivityData.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java @@ -13,27 +13,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.transport.service; +package org.thingsboard.server.common.transport.activity; import lombok.Data; -import org.thingsboard.server.gen.transport.TransportProtos; -/** - * Created by ashvayka on 15.10.18. - */ @Data -public class SessionActivityData { - - private volatile TransportProtos.SessionInfoProto sessionInfo; - private volatile long lastActivityTime; - private volatile long lastReportedActivityTime; - - SessionActivityData(TransportProtos.SessionInfoProto sessionInfo) { - this.sessionInfo = sessionInfo; - } +public class ActivityState { - void updateLastActivityTime() { - this.lastActivityTime = System.currentTimeMillis(); - } + private volatile long lastRecordedTime; + private volatile Metadata metadata; } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java new file mode 100644 index 0000000000..b8fbb9ff12 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +public interface ActivityStrategy { + + boolean onActivity(); + + boolean onReportingPeriodEnd(); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java new file mode 100644 index 0000000000..3f2ec74c3c --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java @@ -0,0 +1,47 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +public enum ActivityStrategyType { + + ALL { + @Override + public ActivityStrategy toStrategy() { + return AllEventsActivityStrategy.getInstance(); + } + }, + FIRST { + @Override + public ActivityStrategy toStrategy() { + return new FirstEventActivityStrategy(); + } + }, + LAST { + @Override + public ActivityStrategy toStrategy() { + return LastEventActivityStrategy.getInstance(); + } + }, + FIRST_AND_LAST { + @Override + public ActivityStrategy toStrategy() { + return new FirstAndLastEventActivityStrategy(); + } + }; + + public abstract ActivityStrategy toStrategy(); + +} 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 new file mode 100644 index 0000000000..f7d2679689 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +public final class AllEventsActivityStrategy implements ActivityStrategy { + + private static final AllEventsActivityStrategy INSTANCE = new AllEventsActivityStrategy(); + + private AllEventsActivityStrategy() { + } + + public static AllEventsActivityStrategy getInstance() { + return INSTANCE; + } + + @Override + public boolean onActivity() { + return true; + } + + @Override + public boolean onReportingPeriodEnd() { + return true; + } + +} 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 new file mode 100644 index 0000000000..a5a59046ae --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode +public final class FirstAndLastEventActivityStrategy implements ActivityStrategy { + + private boolean firstEventReceived; + + @Override + public synchronized boolean onActivity() { + if (!firstEventReceived) { + firstEventReceived = true; + return true; + } + return false; + } + + @Override + public synchronized boolean onReportingPeriodEnd() { + firstEventReceived = false; + return true; + } + +} 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 new file mode 100644 index 0000000000..5353635c45 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import lombok.EqualsAndHashCode; + +@EqualsAndHashCode +public final class FirstEventActivityStrategy implements ActivityStrategy { + + private boolean firstEventReceived; + + @Override + public synchronized boolean onActivity() { + if (!firstEventReceived) { + firstEventReceived = true; + return true; + } + return false; + } + + @Override + public synchronized boolean onReportingPeriodEnd() { + firstEventReceived = false; + return false; + } + +} 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 new file mode 100644 index 0000000000..bce35b0b34 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +public final class LastEventActivityStrategy implements ActivityStrategy { + + private static final LastEventActivityStrategy INSTANCE = new LastEventActivityStrategy(); + + private LastEventActivityStrategy() { + } + + public static LastEventActivityStrategy getInstance() { + return INSTANCE; + } + + @Override + public boolean onActivity() { + return false; + } + + @Override + public boolean onReportingPeriodEnd() { + return true; + } + +} 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 37b77a5268..83928759b0 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 @@ -50,9 +50,9 @@ 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.limit.LimitedApi; +import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; -import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -73,6 +73,11 @@ import org.thingsboard.server.common.transport.TransportResourceCache; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.TransportTenantProfileCache; +import org.thingsboard.server.common.transport.activity.AbstractActivityManager; +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.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; @@ -96,9 +101,9 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.common.AsyncCallbackTemplate; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.provider.TbTransportQueueFactory; import org.thingsboard.server.queue.scheduler.SchedulerComponent; @@ -109,13 +114,11 @@ import org.thingsboard.server.queue.util.TbTransportComponent; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.Collections; -import java.util.HashSet; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Random; -import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -133,13 +136,13 @@ import java.util.stream.Collectors; @Slf4j @Service @TbTransportComponent -public class DefaultTransportService implements TransportService { +public class DefaultTransportService extends AbstractActivityManager implements TransportService { public static final String OVERWRITE_ACTIVITY_TIME = "overwriteActivityTime"; public static final String SESSION_EXPIRED_MESSAGE = "Session has expired due to last activity time!"; public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_OPEN = getSessionEventMsg(TransportProtos.SessionEvent.OPEN); public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_CLOSED = getSessionEventMsg(TransportProtos.SessionEvent.CLOSED); - public static final TransportProtos.SessionCloseNotificationProto SESSION_CLOSE_NOTIFICATION_PROTO = TransportProtos.SessionCloseNotificationProto.newBuilder() + public static final TransportProtos.SessionCloseNotificationProto SESSION_EXPIRED_NOTIFICATION_PROTO = TransportProtos.SessionCloseNotificationProto.newBuilder() .setMessage(SESSION_EXPIRED_MESSAGE).build(); public static final TransportProtos.SubscribeToAttributeUpdatesMsg SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG = TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder() .setSessionType(TransportProtos.SessionType.ASYNC).build(); @@ -156,6 +159,8 @@ public class DefaultTransportService implements TransportService { private long sessionInactivityTimeout; @Value("${transport.sessions.report_timeout}") private long sessionReportTimeout; + @Value("${transport.activity.reporting_strategy:LAST}") + private ActivityStrategyType reportingStrategyType; @Value("${transport.client_side_rpc.timeout:60000}") private long clientSideRpcTimeout; @Value("${queue.transport.poll_interval}") @@ -182,11 +187,9 @@ public class DefaultTransportService implements TransportService { private final TransportRateLimitService rateLimitService; private final DataDecodingEncodingService dataDecodingEncodingService; - private final SchedulerComponent scheduler; private final ApplicationEventPublisher eventPublisher; private final TransportResourceCache transportResourceCache; private final NotificationRuleProcessor notificationRuleProcessor; - private final EntityLimitsCache entityLimitsCache; protected TbQueueRequestTemplate, TbProtoQueueMsg> transportApiRequestTemplate; @@ -202,7 +205,6 @@ public class DefaultTransportService implements TransportService { private ExecutorService mainConsumerExecutor; public final ConcurrentMap sessions = new ConcurrentHashMap<>(); - private final ConcurrentMap sessionsActivity = new ConcurrentHashMap<>(); private final Map toServerRpcPendingMap = new ConcurrentHashMap<>(); private volatile boolean stopped = false; @@ -238,11 +240,11 @@ public class DefaultTransportService implements TransportService { @PostConstruct public void init() { + super.init(); this.ruleEngineProducerStats = statsFactory.createMessagesStats(StatsType.RULE_ENGINE.getName() + ".producer"); this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer"); this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer"); this.transportCallbackExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass()); - this.scheduler.scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); this.scheduler.scheduleAtFixedRate(this::invalidateRateLimits, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); transportApiRequestTemplate = queueProvider.createTransportApiRequestTemplate(); transportApiRequestTemplate.setMessagesStats(transportApiStats); @@ -558,7 +560,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) .setSessionEvent(msg).build(), callback); } @@ -578,7 +580,7 @@ public class DefaultTransportService implements TransportService { } } - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, msg, callback); } } @@ -595,7 +597,7 @@ public class DefaultTransportService implements TransportService { dataPoints += tsKv.getKvCount(); } if (checkLimits(sessionInfo, msg, callback, dataPoints)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); TenantId tenantId = getTenantId(sessionInfo); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); CustomerId customerId = getCustomerId(sessionInfo); @@ -619,7 +621,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback 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 +641,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.GetAttributeRequestMsg msg, TransportServiceCallback 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 +654,7 @@ public class DefaultTransportService implements TransportService { if (sessionMetaData != null) { sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe()); } - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToAttributes(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -665,7 +667,7 @@ public class DefaultTransportService implements TransportService { if (sessionMetaData != null) { sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe()); } - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToRPC(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); } @@ -674,7 +676,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcResponseMsg msg, TransportServiceCallback 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 +685,7 @@ public class DefaultTransportService implements TransportService { @Override public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback); } } @@ -704,7 +706,7 @@ public class DefaultTransportService implements TransportService { if (checkLimits(sessionInfo, responseMsg, callback)) { if (reportActivity) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); } sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY)); @@ -736,7 +738,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToServerRpcRequestMsg msg, TransportServiceCallback callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); UUID sessionId = toSessionId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo); DeviceId deviceId = getDeviceId(sessionInfo); @@ -761,7 +763,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ClaimDeviceMsg msg, TransportServiceCallback callback) { if (checkLimits(sessionInfo, msg, callback)) { - reportActivityInternal(sessionInfo); + recordActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) .setClaimDevice(msg).build(), callback); } @@ -780,69 +782,106 @@ public class DefaultTransportService implements TransportService { } @Override - public void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { - reportActivityInternal(sessionInfo); + public void recordActivity(TransportProtos.SessionInfoProto sessionInfo) { + recordActivityInternal(sessionInfo); } - private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { - UUID sessionId = toSessionId(sessionInfo); - SessionActivityData sessionMetaData = sessionsActivity.computeIfAbsent(sessionId, id -> new SessionActivityData(sessionInfo)); - sessionMetaData.updateLastActivityTime(); - } - - private void checkInactivityAndReportActivity() { - long expTime = System.currentTimeMillis() - sessionInactivityTimeout; - Set sessionsToRemove = new HashSet<>(); - sessionsActivity.forEach((uuid, sessionAD) -> { - long lastActivityTime = sessionAD.getLastActivityTime(); - SessionMetaData sessionMD = sessions.get(uuid); - if (sessionMD != null) { - sessionAD.setSessionInfo(sessionMD.getSessionInfo()); - } else { - sessionsToRemove.add(uuid); - } - TransportProtos.SessionInfoProto sessionInfo = sessionAD.getSessionInfo(); - - if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) { - var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB()); - SessionMetaData gwMetaData = sessions.get(gwSessionId); - SessionActivityData gwActivityData = sessionsActivity.get(gwSessionId); - if (gwMetaData != null && gwMetaData.isOverwriteActivityTime()) { - lastActivityTime = Math.max(gwActivityData.getLastActivityTime(), lastActivityTime); - } + private void recordActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { + onActivity(toSessionId(sessionInfo), getCurrentTimeMillis()); + } + + @Override + protected long getReportingPeriodMillis() { + return sessionReportTimeout; + } + + @Override + protected ActivityState createNewState(UUID sessionId) { + SessionMetaData session = sessions.get(sessionId); + if (session == null) { + return null; + } + ActivityState state = new ActivityState<>(); + state.setMetadata(session.getSessionInfo()); + return state; + } + + @Override + protected ActivityStrategy getStrategy() { + return reportingStrategyType.toStrategy(); + } + + @Override + protected ActivityState updateState(UUID sessionId, ActivityState state) { + SessionMetaData session = sessions.get(sessionId); + if (session == null) { + return null; + } + + state.setMetadata(session.getSessionInfo()); + var sessionInfo = state.getMetadata(); + + if (sessionInfo.getGwSessionIdMSB() == 0L || sessionInfo.getGwSessionIdLSB() == 0L) { + return state; + } + + var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB()); + SessionMetaData gwSession = sessions.get(gwSessionId); + if (gwSession == null || !gwSession.isOverwriteActivityTime()) { + return state; + } + + long lastRecordedTime = state.getLastRecordedTime(); + long gwLastRecordedTime = getLastRecordedTime(gwSessionId); + log.debug("Session with id: [{}] has gateway session with id: [{}] with overwrite activity time enabled. " + + "Updating last activity time. Session last recorded time: [{}], gateway session last recorded time: [{}].", + sessionId, gwSessionId, lastRecordedTime, gwLastRecordedTime); + state.setLastRecordedTime(Math.max(lastRecordedTime, gwLastRecordedTime)); + return state; + } + + @Override + protected boolean hasExpired(long lastRecordedTime) { + return (getCurrentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime; + } + + long getCurrentTimeMillis() { + return System.currentTimeMillis(); + } + + @Override + protected void onStateExpiry(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) { + log.debug("Session with id: [{}] has expired due to last activity time.", sessionId); + SessionMetaData expiredSession = sessions.remove(sessionId); + if (expiredSession != null) { + deregisterSession(sessionInfo); + process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); + expiredSession.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); + } + } + + @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: [{}].", 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); + } - if (lastActivityTime < expTime) { - if (sessionMD != null) { - if (log.isDebugEnabled()) { - log.debug("[{}] Session has expired due to last activity time: {}", toSessionId(sessionInfo), lastActivityTime); - } - sessions.remove(uuid); - sessionsToRemove.add(uuid); - process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); - sessionMD.getListener().onRemoteSessionCloseCommand(uuid, SESSION_CLOSE_NOTIFICATION_PROTO); - } - } else { - if (lastActivityTime > sessionAD.getLastReportedActivityTime()) { - final long lastActivityTimeFinal = lastActivityTime; - process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder() - .setAttributeSubscription(sessionMD != null && sessionMD.isSubscribedToAttributes()) - .setRpcSubscription(sessionMD != null && sessionMD.isSubscribedToRPC()) - .setLastActivityTime(lastActivityTime).build(), new TransportServiceCallback() { - @Override - public void onSuccess(Void msg) { - sessionAD.setLastReportedActivityTime(lastActivityTimeFinal); - } - @Override - public void onError(Throwable e) { - log.warn("[{}] Failed to report last activity time", uuid, e); - } - }); - } + @Override + public void onError(Throwable e) { + callback.onFailure(sessionId, e); } }); - // Removes all closed or short-lived sessions. - sessionsToRemove.forEach(sessionsActivity::remove); } @Override @@ -913,7 +952,7 @@ public class DefaultTransportService implements TransportService { } TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); - builder.setTs(TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) * 1000L + (atomicTs.getAndIncrement() % 1000)); + builder.setTs(TimeUnit.MILLISECONDS.toSeconds(getCurrentTimeMillis()) * 1000L + (atomicTs.getAndIncrement() % 1000)); builder.addKv(TransportProtos.KeyValueProto.newBuilder() .setKey("transportLog") .setType(TransportProtos.KeyValueType.STRING_V) diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyTypeTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyTypeTest.java new file mode 100644 index 0000000000..7dcf0e8fd2 --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyTypeTest.java @@ -0,0 +1,44 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +public class ActivityStrategyTypeTest { + + @Test + public void testCreateAllEventsStrategy() { + assertThat(ActivityStrategyType.ALL.toStrategy()).isEqualTo(AllEventsActivityStrategy.getInstance()); + } + + @Test + public void testCreateFirstEventStrategy() { + assertThat(ActivityStrategyType.FIRST.toStrategy()).isEqualTo(new FirstEventActivityStrategy()); + } + + @Test + public void testCreateLastEventStrategy() { + assertThat(ActivityStrategyType.LAST.toStrategy()).isEqualTo(LastEventActivityStrategy.getInstance()); + } + + @Test + public void testCreateFirstAndLastEventStrategy() { + assertThat(ActivityStrategyType.FIRST_AND_LAST.toStrategy()).isEqualTo(new FirstAndLastEventActivityStrategy()); + } + +} diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategyTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategyTest.java new file mode 100644 index 0000000000..ea89b28a02 --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategyTest.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class AllEventsActivityStrategyTest { + + private AllEventsActivityStrategy strategy; + + @BeforeEach + public void setUp() { + strategy = AllEventsActivityStrategy.getInstance(); + } + + @Test + public void testOnActivity() { + assertTrue(strategy.onActivity(), "onActivity() should always return true."); + } + + @Test + public void testOnReportingPeriodEnd() { + assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); + } + +} diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategyTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategyTest.java new file mode 100644 index 0000000000..9ca6b58502 --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategyTest.java @@ -0,0 +1,53 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class FirstAndLastEventActivityStrategyTest { + + private FirstAndLastEventActivityStrategy strategy; + + @BeforeEach + public void setUp() { + strategy = new FirstAndLastEventActivityStrategy(); + } + + @Test + public void testOnActivity_FirstCall() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + } + + @Test + public void testOnActivity_SubsequentCalls() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false."); + } + + @Test + public void testOnReportingPeriodEnd() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); + assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd() for the next reporting period"); + assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); + } + + +} diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategyTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategyTest.java new file mode 100644 index 0000000000..0df3517098 --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategyTest.java @@ -0,0 +1,52 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class FirstEventActivityStrategyTest { + + private FirstEventActivityStrategy strategy; + + @BeforeEach + public void setUp() { + strategy = new FirstEventActivityStrategy(); + } + + @Test + public void testOnActivity_FirstCall() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + } + + @Test + public void testOnActivity_SubsequentCalls() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false."); + } + + @Test + public void testOnReportingPeriodEnd() { + assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); + assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false."); + assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd()."); + assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false."); + } + +} diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategyTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategyTest.java new file mode 100644 index 0000000000..d747c13787 --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategyTest.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class LastEventActivityStrategyTest { + + private LastEventActivityStrategy strategy; + + @BeforeEach + public void setUp() { + strategy = LastEventActivityStrategy.getInstance(); + } + + @Test + public void testOnActivity() { + assertFalse(strategy.onActivity(), "onActivity() should always return false."); + } + + @Test + public void testOnReportingPeriodEnd() { + assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); + } + +} diff --git a/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java new file mode 100644 index 0000000000..e0f5ef128e --- /dev/null +++ b/common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java @@ -0,0 +1,508 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.service; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.MethodSource; +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.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.gen.transport.TransportProtos; + +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.stream.Stream; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.doCallRealMethod; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_CLOSED; +import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EXPIRED_NOTIFICATION_PROTO; + +@ExtendWith(MockitoExtension.class) +public class TransportActivityManagerTest { + + private final UUID SESSION_ID = UUID.fromString("1306648a-9b26-11ee-b9d1-0242ac120002"); + + @Mock + private DefaultTransportService transportServiceMock; + private ConcurrentMap sessions; + + @BeforeEach + public void setup() { + sessions = new ConcurrentHashMap<>(); + ReflectionTestUtils.setField(transportServiceMock, "sessions", sessions); + } + + @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); + long expectedTime = 123L; + when(transportServiceMock.getCurrentTimeMillis()).thenReturn(expectedTime); + + // WHEN + transportServiceMock.recordActivity(sessionInfo); + + // THEN + verify(transportServiceMock).onActivity(SESSION_ID, expectedTime); + } + + @Test + void givenKey_whenCreatingNewState_thenShouldCorrectlyCreateNewEmptyState() { + // GIVEN + TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() + .setSessionIdMSB(SESSION_ID.getMostSignificantBits()) + .setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) + .build(); + sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, null)); + + when(transportServiceMock.createNewState(SESSION_ID)).thenCallRealMethod(); + + ActivityState expectedState = new ActivityState<>(); + expectedState.setMetadata(sessionInfo); + + // WHEN + ActivityState actualState = transportServiceMock.createNewState(SESSION_ID); + + // THEN + assertThat(actualState).isEqualTo(expectedState); + } + + @ParameterizedTest + @EnumSource(ActivityStrategyType.class) + void givenDifferentReportingStrategies_whenGettingStrategy_thenShouldReturnCorrectStrategy(ActivityStrategyType reportingStrategyType) { + // GIVEN + doCallRealMethod().when(transportServiceMock).getStrategy(); + ReflectionTestUtils.setField(transportServiceMock, "reportingStrategyType", reportingStrategyType); + + // WHEN + ActivityStrategy actualStrategy = transportServiceMock.getStrategy(); + + // THEN + assertThat(actualStrategy).isEqualTo(reportingStrategyType.toStrategy()); + } + + @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.setMetadata(sessionInfo); + + 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; + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + + 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.getMetadata()).isEqualTo(sessionInfo); + } + + @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; + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + + 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.getMetadata()).isEqualTo(sessionInfo); + + 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; + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + + 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.getMetadata()).isEqualTo(sessionInfo); + + 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; + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + + 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.getMetadata()).isEqualTo(sessionInfo); + } + + @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; + + ActivityState state = new ActivityState<>(); + state.setLastRecordedTime(lastRecordedTime); + state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); + + 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.getMetadata()).isEqualTo(sessionInfo); + } + + @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); + } + +}