From 22073936e09a796178fa3289e927b8fa6cef4140 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Fri, 8 Dec 2023 12:32:13 +0200 Subject: [PATCH] Add logging --- ...irstAndLastIntegrationActivityManager.java | 13 +++++++--- .../FirstOnlyIntegrationActivityManager.java | 13 +++++++--- .../LastOnlyIntegrationActivityManager.java | 7 +++-- .../activity/AbstractActivityManager.java | 26 ++++++++++++------- .../service/DefaultTransportService.java | 4 ++- .../FirstAndLastTransportActivityManager.java | 14 +++++----- .../FirstOnlyTransportActivityManager.java | 14 +++++----- .../LastOnlyTransportActivityManager.java | 11 ++++---- 8 files changed, 67 insertions(+), 35 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java index e0dcc97ac0..c9d05eaa26 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java @@ -81,6 +81,8 @@ public class FirstAndLastIntegrationActivityManager extends AbstractActivityMana return activityStateWrapper; } if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { + log.debug("[{}][{}] Going to report first activity event for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { @Override public void onSuccess(IntegrationActivityKey key, long reportedTime) { @@ -104,13 +106,14 @@ public class FirstAndLastIntegrationActivityManager extends AbstractActivityMana @Override public void onFailure(@NonNull Throwable t) { - log.debug("[{}] Failed to report first activity event in a period for device with id: [{}].", activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); + log.debug("[{}][{}] Failed to report first activity event for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); } }, MoreExecutors.directExecutor()); } @Override - protected void onReportingPeriodEnd() { + protected void doOnReportingPeriodEnd() { for (Map.Entry entry : states.entrySet()) { var activityKey = entry.getKey(); var activityStateWrapper = entry.getValue(); @@ -118,9 +121,12 @@ public class FirstAndLastIntegrationActivityManager extends AbstractActivityMana long lastRecordedTime = activityState.getLastRecordedTime(); // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks if (!activityStateWrapper.isAlreadyBeenReported()) { + log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); states.remove(activityKey); } if (activityState.getLastReportedTime() < lastRecordedTime) { + log.debug("[{}][{}] Going to report last activity event for device with id: [{}].", activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); reporter.report(activityKey, lastRecordedTime, activityState, new ActivityReportCallback<>() { @Override public void onSuccess(IntegrationActivityKey key, long newLastReportedTime) { @@ -129,7 +135,8 @@ public class FirstAndLastIntegrationActivityManager extends AbstractActivityMana @Override public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}] Failed to report last activity event in a period for device with id: [{}].", activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); + log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); } }); } diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java index 38f14f538a..35567847a5 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java @@ -81,6 +81,8 @@ public class FirstOnlyIntegrationActivityManager extends AbstractActivityManager return activityStateWrapper; } if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { + log.debug("[{}][{}] Going to report first activity event for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { @Override public void onSuccess(IntegrationActivityKey key, long reportedTime) { @@ -104,13 +106,14 @@ public class FirstOnlyIntegrationActivityManager extends AbstractActivityManager @Override public void onFailure(@NonNull Throwable t) { - log.debug("[{}] Failed to report first activity event in a period for device with id: [{}].", activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); + log.debug("[{}][{}] Failed to report activity event for device with id: [{}].", + name, activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); } }, MoreExecutors.directExecutor()); } @Override - protected void onReportingPeriodEnd() { + protected void doOnReportingPeriodEnd() { for (Map.Entry entry : states.entrySet()) { var activityKey = entry.getKey(); var activityStateWrapper = entry.getValue(); @@ -118,9 +121,12 @@ public class FirstOnlyIntegrationActivityManager extends AbstractActivityManager long lastRecordedTime = activityState.getLastRecordedTime(); // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks if (!activityStateWrapper.isAlreadyBeenReported()) { + log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); states.remove(activityKey); // report leftover events if (activityState.getLastReportedTime() < lastRecordedTime) { + log.debug("[{}][{}] Going to report leftover activity event for device with id: [{}].", activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); reporter.report(activityKey, lastRecordedTime, activityState, new ActivityReportCallback<>() { @Override public void onSuccess(IntegrationActivityKey key, long reportedTime) { @@ -129,7 +135,8 @@ public class FirstOnlyIntegrationActivityManager extends AbstractActivityManager @Override public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}] Failed to report last activity event in a period for device with id: [{}].", activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); + log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); } }); } diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java index ef6fa4f8b6..6517a1af2a 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java @@ -64,7 +64,7 @@ public class LastOnlyIntegrationActivityManager extends AbstractActivityManager< } @Override - public void onReportingPeriodEnd() { + public void doOnReportingPeriodEnd() { long expirationTime = System.currentTimeMillis() - reportingPeriodMillis; for (Map.Entry entry : states.entrySet()) { var activityKey = entry.getKey(); @@ -72,6 +72,8 @@ public class LastOnlyIntegrationActivityManager extends AbstractActivityManager< long lastRecordedTime = activityState.getLastRecordedTime(); // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks if (lastRecordedTime < expirationTime) { + log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); states.remove(activityKey); } if (activityState.getLastReportedTime() < lastRecordedTime) { @@ -83,7 +85,8 @@ public class LastOnlyIntegrationActivityManager extends AbstractActivityManager< @Override public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}] Failed to report last activity event in a period for device with id: [{}].", activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); + log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", + activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); } }); } 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 4dc84037a8..9d7b5c4e19 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 @@ -44,7 +44,7 @@ import java.util.function.Supplier; @Slf4j public abstract class AbstractActivityManager implements ActivityManager { - private String name; + protected String name; protected long reportingPeriodMillis; protected ActivityStateReporter reporter; private ScheduledExecutorService scheduler; @@ -52,17 +52,19 @@ public abstract class AbstractActivityManager @Override public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter reporter) { - this.name = StringUtils.notBlankOrDefault(name, "activity-manager"); - this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided activity reporter is null."); - if (reportingPeriodMillis <= 0) { - reportingPeriodMillis = 3000; - log.error("[{}] Negative or zero reporting period millisecond was provided. Going to use reporting period value of 3 seconds.", this.name); - } - this.reportingPeriodMillis = reportingPeriodMillis; if (!initialized) { + log.info("[{}] initializing.", this.name); + this.name = StringUtils.notBlankOrDefault(name, "activity-manager"); + this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided activity reporter is null."); + if (reportingPeriodMillis <= 0) { + reportingPeriodMillis = 3000; + log.error("[{}] Negative or zero reporting period millisecond was provided. Going to use reporting period value of 3 seconds.", this.name); + } + this.reportingPeriodMillis = reportingPeriodMillis; scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(this.name)); scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodMillis), reportingPeriodMillis, TimeUnit.MILLISECONDS); initialized = true; + log.info("[{}] initialized.", this.name); } } @@ -80,12 +82,18 @@ public abstract class AbstractActivityManager log.error("[{}] Failed to process activity event: provided new activity state supplier is null.", name); return; } + log.debug("[{}] Received activity event for key: [{}]", name, key); doOnActivity(key, newStateSupplier); } protected abstract void doOnActivity(Key key, Supplier newStateSupplier); - protected abstract void onReportingPeriodEnd(); + private void onReportingPeriodEnd() { + log.debug("[{}] Going to end reporting period.", name); + doOnReportingPeriodEnd(); + } + + protected abstract void doOnReportingPeriodEnd(); @Override public synchronized void destroy() { 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 692421b0ab..d2fa8ad8ca 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 @@ -187,6 +187,7 @@ public class DefaultTransportService implements TransportService { private final NotificationRuleProcessor notificationRuleProcessor; private final EntityLimitsCache entityLimitsCache; private ActivityManager activityManager; + private static final String ACTIVITY_MANAGER_NAME = "transport-activity-manager"; protected TbQueueRequestTemplate, TbProtoQueueMsg> transportApiRequestTemplate; protected TbQueueProducer> ruleEngineMsgProducer; @@ -255,10 +256,11 @@ public class DefaultTransportService implements TransportService { transportNotificationsConsumer.subscribe(Collections.singleton(tpi)); transportApiRequestTemplate.init(); mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); - activityManager.init("transport-activity-manager", sessionReportTimeout, activityReporter); + activityManager.init(ACTIVITY_MANAGER_NAME, sessionReportTimeout, activityReporter); } private final ActivityStateReporter activityReporter = (sessionId, timeToReport, state, reportCallback) -> { + log.debug("[{}] Reporting activity state for key: [{}]. Time to report: [{}].", ACTIVITY_MANAGER_NAME, sessionId, timeToReport); SessionMetaData sessionMetaData = sessions.get(sessionId); TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() .setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes()) diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java index 6370272813..db0ebc7ddb 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java @@ -98,6 +98,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage return activityStateWrapper; } if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { + log.debug("[{}] Going to report first activity event for session with id: [{}]", name, sessionId); reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { @Override public void onSuccess(UUID key, long reportedTime) { @@ -121,13 +122,13 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage @Override public void onFailure(@NonNull Throwable t) { - log.debug("[{}] Failed to report first activity event in a period.", sessionId); + log.debug("[{}] Failed to report first activity event for session with id: [{}]", name, sessionId); } }, MoreExecutors.directExecutor()); } @Override - protected void onReportingPeriodEnd() { + protected void doOnReportingPeriodEnd() { Set statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); @@ -138,6 +139,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage if (sessionMetaData != null) { activityState.setSessionInfoProto(sessionMetaData.getSessionInfo()); } else { + log.debug("[{}] Session with id: [{}] is not present. Marking it's activity state for removal.", name, sessionId); statesToRemove.add(sessionId); } @@ -150,6 +152,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage if (gwSessionMetaData != null && gwSessionMetaData.isOverwriteActivityTime()) { ActivityStateWrapper gwActivityStateWrapper = states.get(gwSessionId); if (gwActivityStateWrapper != null) { + log.debug("[{}] Session with id: [{}] has gateway session with id: [{}] with overwrite activity time enabled. Updating last activity time.", name, sessionId, gwSessionId); lastActivityTime = Math.max(gwActivityStateWrapper.getState().getLastRecordedTime(), lastActivityTime); } } @@ -158,15 +161,14 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; boolean hasExpired = sessionMetaData != null && lastActivityTime < expirationTime; if (hasExpired) { - if (log.isDebugEnabled()) { - log.debug("[{}] Session has expired due to last activity time: [{}].", sessionId, lastActivityTime); - } + log.debug("[{}] Session with id: [{}] has expired due to last activity time: [{}]. Marking it's activity state for removal.", name, sessionId, lastActivityTime); statesToRemove.add(sessionId); transportService.deregisterSession(sessionInfo); transportService.process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); } if (activityState.getLastReportedTime() < lastActivityTime) { + log.debug("[{}] Going to report last activity event for session with id: [{}].", name, sessionId); reporter.report(sessionId, lastActivityTime, activityState, new ActivityReportCallback<>() { @Override public void onSuccess(UUID key, long reportedTime) { @@ -175,7 +177,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage @Override public void onFailure(UUID key, Throwable t) { - log.debug("[{}] Failed to report last activity event in a period.", sessionId); + log.debug("[{}] Failed to report last activity event for session with id: [{}].", name, sessionId); } }); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java index 768796b2be..b0762e0b5e 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java @@ -98,6 +98,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager() { @Override public void onSuccess(UUID key, long reportedTime) { @@ -121,13 +122,13 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); @@ -138,6 +139,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager() { @Override public void onSuccess(UUID key, long reportedTime) { @@ -176,7 +178,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); @@ -91,6 +91,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager() { @Override public void onSuccess(UUID key, long reportedTime) { @@ -128,7 +129,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager