Browse Source

Add logging

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
22073936e0
  1. 13
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java
  2. 13
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java
  3. 7
      application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java
  4. 26
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  5. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  6. 14
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java
  7. 14
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java
  8. 11
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java

13
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<IntegrationActivityKey, ActivityStateWrapper> 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());
}
});
}

13
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<IntegrationActivityKey, ActivityStateWrapper> 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());
}
});
}

7
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<IntegrationActivityKey, ActivityState> 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());
}
});
}

26
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<Key, State extends ActivityState> implements ActivityManager<Key, State> {
private String name;
protected String name;
protected long reportingPeriodMillis;
protected ActivityStateReporter<Key, State> reporter;
private ScheduledExecutorService scheduler;
@ -52,17 +52,19 @@ public abstract class AbstractActivityManager<Key, State extends ActivityState>
@Override
public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, State> 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<Key, State extends ActivityState>
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<State> 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() {

4
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<UUID, TransportActivityState> activityManager;
private static final String ACTIVITY_MANAGER_NAME = "transport-activity-manager";
protected TbQueueRequestTemplate<TbProtoQueueMsg<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate;
protected TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> 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<UUID, TransportActivityState> 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())

14
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<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, ActivityStateWrapper> 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);
}
});
}

14
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<U
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 FirstOnlyTransportActivityManager extends AbstractActivityManager<U
@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<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, ActivityStateWrapper> entry : states.entrySet()) {
var sessionId = entry.getKey();
@ -138,6 +139,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
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 FirstOnlyTransportActivityManager extends AbstractActivityManager<U
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,9 +161,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
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);
@ -168,6 +169,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
}
boolean shouldReportLeftoverEvents = sessionMetaData == null || hasExpired;
if (shouldReportLeftoverEvents && activityState.getLastReportedTime() < lastActivityTime) {
log.debug("[{}] Going to report leftover activity event for session with id: [{}].", name, sessionId);
reporter.report(sessionId, lastActivityTime, activityState, new ActivityReportCallback<>() {
@Override
public void onSuccess(UUID key, long reportedTime) {
@ -176,7 +178,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
@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 leftover activity event for session with id: [{}].", name, sessionId);
}
});
}

11
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java

@ -81,7 +81,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
}
@Override
protected void onReportingPeriodEnd() {
protected void doOnReportingPeriodEnd() {
Set<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, TransportActivityState> entry : states.entrySet()) {
var sessionId = entry.getKey();
@ -91,6 +91,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
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);
}
@ -103,6 +104,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
if (gwSessionMetaData != null && gwSessionMetaData.isOverwriteActivityTime()) {
TransportActivityState gwActivityState = states.get(gwSessionId);
if (gwActivityState != 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(gwActivityState.getLastRecordedTime(), lastActivityTime);
}
}
@ -111,15 +113,14 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
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) {
@ -128,7 +129,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
@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);
}
});
}

Loading…
Cancel
Save