Browse Source

Monitoring service to single-thread

pull/7725/head
ViacheslavKlimov 3 years ago
parent
commit
869136da89
  1. 70
      monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java
  2. 72
      monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java
  3. 2
      monitoring/src/main/java/org/thingsboard/monitoring/client/WsClientFactory.java
  4. 17
      monitoring/src/main/java/org/thingsboard/monitoring/config/TransportType.java
  5. 2
      monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringConfig.java
  6. 2
      monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringConfig.java
  7. 2
      monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringConfig.java
  8. 4
      monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringConfig.java
  9. 2
      monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/CmdsWrapper.java
  10. 15
      monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataCmd.java
  11. 44
      monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java
  12. 28
      monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/LatestValueCmd.java
  13. 52
      monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java
  14. 74
      monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java
  15. 65
      monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportHealthChecker.java
  16. 143
      monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportMonitoringService.java
  17. 10
      monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/CoapTransportHealthChecker.java
  18. 10
      monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/HttpTransportHealthChecker.java
  19. 10
      monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/MqttTransportHealthChecker.java
  20. 2
      monitoring/src/main/resources/logback.xml
  21. 11
      monitoring/src/main/resources/tb-monitoring.yml

70
monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java

@ -16,31 +16,13 @@
package org.thingsboard.monitoring;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationRunner;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.monitoring.client.TbClient;
import org.thingsboard.monitoring.config.DeviceConfig;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.service.TransportMonitoringServiceConfig;
import org.thingsboard.monitoring.service.TransportMonitoringService;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration;
import org.thingsboard.server.common.data.device.data.DefaultDeviceTransportConfiguration;
import org.thingsboard.server.common.data.device.data.DeviceData;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
@ -55,56 +37,4 @@ public class ThingsboardMonitoringApplication {
.run(args);
}
@Bean
public ApplicationRunner initAndStartMonitoringServices(List<TransportMonitoringServiceConfig> configs, TbClient tbClient, ApplicationContext context) {
return args -> {
List<TransportMonitoringService<?>> monitoringServices = new LinkedList<>();
configs.forEach(config -> {
config.getTargets().stream()
.filter(target -> StringUtils.isNotBlank(target.getBaseUrl()))
.peek(target -> checkMonitoringTarget(config, target, tbClient))
.forEach(target -> {
TransportMonitoringService<?> monitoringService = context.getBean(config.getTransportType().getMonitoringServiceClass(), config, target);
monitoringServices.add(monitoringService);
});
});
monitoringServices.forEach(TransportMonitoringService::startMonitoring);
};
}
private void checkMonitoringTarget(TransportMonitoringServiceConfig config, MonitoringTargetConfig target, TbClient tbClient) {
DeviceConfig deviceConfig = target.getDevice();
tbClient.logIn();
DeviceId deviceId;
if (deviceConfig == null || deviceConfig.getId() == null) {
String deviceName = String.format("[%s] Monitoring device (%s)", config.getTransportType(), target.getBaseUrl());
Device device = tbClient.getTenantDevice(deviceName)
.orElseGet(() -> {
log.info("Creating new device '{}'", deviceName);
Device monitoringDevice = new Device();
monitoringDevice.setName(deviceName);
monitoringDevice.setType("default");
DeviceData deviceData = new DeviceData();
deviceData.setConfiguration(new DefaultDeviceConfiguration());
deviceData.setTransportConfiguration(new DefaultDeviceTransportConfiguration());
return tbClient.saveDevice(monitoringDevice);
});
deviceId = device.getId();
target.getDevice().setId(deviceId.toString());
} else {
deviceId = new DeviceId(deviceConfig.getId());
}
log.debug("Loading credentials for device {}", deviceId);
DeviceCredentials credentials = tbClient.getDeviceCredentialsByDeviceId(deviceId)
.orElseThrow(() -> new IllegalArgumentException("No credentials found for device " + deviceId));
target.getDevice().setCredentials(credentials);
}
@Bean
public ScheduledExecutorService monitoringExecutor(@Value("${monitoring.monitoring_executor_thread_pool_size}") int threadPoolSize) {
return Executors.newScheduledThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("monitoring-executor"));
}
}

72
monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java

@ -23,30 +23,42 @@ import org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.monitoring.data.cmd.CmdsWrapper;
import org.thingsboard.monitoring.data.cmd.TimeseriesSubscriptionCmd;
import org.thingsboard.monitoring.data.cmd.TimeseriesUpdate;
import org.thingsboard.monitoring.data.cmd.EntityDataCmd;
import org.thingsboard.monitoring.data.cmd.EntityDataUpdate;
import org.thingsboard.monitoring.data.cmd.LatestValueCmd;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.query.EntityDataPageLink;
import org.thingsboard.server.common.data.query.EntityDataQuery;
import org.thingsboard.server.common.data.query.EntityKey;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.EntityListFilter;
import javax.net.ssl.SSLParameters;
import java.net.URI;
import java.nio.channels.NotYetConnectedException;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
@Slf4j
public class WsClient extends WebSocketClient implements AutoCloseable {
public volatile String lastMsg;
public volatile JsonNode lastMsg;
private CountDownLatch reply;
private CountDownLatch update;
private final Lock updateLock = new ReentrantLock();
public WsClient(URI serverUri) {
private long requestTimeoutMs;
public WsClient(URI serverUri, long requestTimeoutMs) {
super(serverUri);
this.requestTimeoutMs = requestTimeoutMs;
}
@Override
@ -56,13 +68,13 @@ public class WsClient extends WebSocketClient implements AutoCloseable {
@Override
public void onMessage(String s) {
log.trace("Received new msg: {}", s);
if (s == null) {
return;
}
updateLock.lock();
try {
lastMsg = s;
lastMsg = JacksonUtil.toJsonNode(s);
log.trace("Received new msg: {}", lastMsg.toPrettyString());
if (update != null) {
update.countDown();
}
@ -106,18 +118,25 @@ public class WsClient extends WebSocketClient implements AutoCloseable {
super.send(text);
}
public void subscribeForTelemetry(UUID deviceId, String telemetryKey) {
TimeseriesSubscriptionCmd subCmd = new TimeseriesSubscriptionCmd();
subCmd.setEntityType("DEVICE");
subCmd.setEntityId(deviceId.toString());
subCmd.setScope("LATEST_TELEMETRY");
subCmd.setKeys(telemetryKey);
subCmd.setCmdId(RandomUtils.nextInt(0, 100));
public WsClient subscribeForTelemetry(List<UUID> devices, String key) {
EntityDataCmd cmd = new EntityDataCmd();
cmd.setCmdId(RandomUtils.nextInt(0, 1000));
EntityListFilter devicesFilter = new EntityListFilter();
devicesFilter.setEntityType(EntityType.DEVICE);
devicesFilter.setEntityList(devices.stream().map(UUID::toString).collect(Collectors.toList()));
EntityDataPageLink pageLink = new EntityDataPageLink(100,0, null, null);
EntityDataQuery devicesQuery = new EntityDataQuery(devicesFilter, pageLink, Collections.emptyList(), Collections.emptyList(), Collections.emptyList());
cmd.setQuery(devicesQuery);
LatestValueCmd latestCmd = new LatestValueCmd();
latestCmd.setKeys(List.of(new EntityKey(EntityKeyType.TIME_SERIES, key)));
cmd.setLatestCmd(latestCmd);
CmdsWrapper wrapper = new CmdsWrapper();
wrapper.setTsSubCmds(List.of(subCmd));
wrapper.setEntityDataCmds(List.of(cmd));
send(JacksonUtil.toString(wrapper));
log.trace("Subscribed for telemetry (key: {})", telemetryKey);
return this;
}
public JsonNode waitForUpdate(long ms) {
@ -134,38 +153,37 @@ public class WsClient extends WebSocketClient implements AutoCloseable {
return null;
}
public JsonNode waitForReply(int ms) {
public JsonNode waitForReply() {
try {
if (reply.await(ms, TimeUnit.MILLISECONDS)) {
if (reply.await(requestTimeoutMs, TimeUnit.MILLISECONDS)) {
log.trace("Waited for reply");
return getLastMsg();
}
} catch (InterruptedException e) {
log.debug("Failed to await reply", e);
}
log.trace("No reply arrived within {} ms", ms);
return null;
log.trace("No reply arrived within {} ms", requestTimeoutMs);
throw new IllegalStateException("No WS reply arrived within " + requestTimeoutMs + " ms");
}
public JsonNode getLastMsg() {
JsonNode msg = JacksonUtil.toJsonNode(lastMsg);
if (msg != null) {
JsonNode errorMsg = msg.get("errorMsg");
private JsonNode getLastMsg() {
if (lastMsg != null) {
JsonNode errorMsg = lastMsg.get("errorMsg");
if (errorMsg != null && !errorMsg.isNull() && StringUtils.isNotEmpty(errorMsg.asText())) {
throw new RuntimeException("WS error from server: " + errorMsg.asText());
} else {
return msg;
return lastMsg;
}
} else {
return null;
}
}
public Object getTelemetryKeyUpdate(String key) {
public Object getTelemetryUpdate(UUID deviceId, String key) {
JsonNode lastMsg = getLastMsg();
if (lastMsg == null || lastMsg.isNull()) return null;
TimeseriesUpdate update = JacksonUtil.treeToValue(lastMsg, TimeseriesUpdate.class);
return update.getLatest(key);
EntityDataUpdate update = JacksonUtil.treeToValue(lastMsg, EntityDataUpdate.class);
return update.getLatest(deviceId, key);
}
@Override

2
monitoring/src/main/java/org/thingsboard/monitoring/client/WsClientFactory.java

@ -39,7 +39,7 @@ public class WsClientFactory {
public WsClient createClient(String accessToken) throws Exception {
URI uri = new URI(wsConfig.getBaseUrl() + "/api/ws/plugins/telemetry?token=" + accessToken);
stopWatch.start();
WsClient wsClient = new WsClient(uri);
WsClient wsClient = new WsClient(uri, wsConfig.getRequestTimeoutMs());
if (wsConfig.getBaseUrl().startsWith("wss")) {
SSLContextBuilder builder = SSLContexts.custom();
builder.loadTrustMaterial(null, (TrustStrategy) (chain, authType) -> true);

17
monitoring/src/main/java/org/thingsboard/monitoring/config/TransportType.java

@ -17,18 +17,19 @@ package org.thingsboard.monitoring.config;
import lombok.AllArgsConstructor;
import lombok.Getter;
import org.thingsboard.monitoring.service.TransportMonitoringService;
import org.thingsboard.monitoring.service.impl.CoapTransportMonitoringService;
import org.thingsboard.monitoring.service.impl.HttpTransportMonitoringService;
import org.thingsboard.monitoring.service.impl.MqttTransportMonitoringService;
import org.thingsboard.monitoring.transport.TransportHealthChecker;
import org.thingsboard.monitoring.transport.impl.CoapTransportHealthChecker;
import org.thingsboard.monitoring.transport.impl.HttpTransportHealthChecker;
import org.thingsboard.monitoring.transport.impl.MqttTransportHealthChecker;
@AllArgsConstructor
@Getter
public enum TransportType {
MQTT(MqttTransportMonitoringService.class),
COAP(CoapTransportMonitoringService.class),
HTTP(HttpTransportMonitoringService.class);
private final Class<? extends TransportMonitoringService<?>> monitoringServiceClass;
MQTT(MqttTransportHealthChecker.class),
COAP(CoapTransportHealthChecker.class),
HTTP(HttpTransportHealthChecker.class);
private final Class<? extends TransportHealthChecker<?>> serviceClass;
}

2
monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringServiceConfig.java → monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringConfig.java

@ -23,7 +23,7 @@ import org.thingsboard.monitoring.config.TransportType;
@Component
@ConditionalOnProperty(name = "monitoring.transports.coap.enabled", havingValue = "true")
@ConfigurationProperties(prefix = "monitoring.transports.coap")
public class CoapTransportMonitoringServiceConfig extends TransportMonitoringServiceConfig {
public class CoapTransportMonitoringConfig extends TransportMonitoringConfig {
@Override
public TransportType getTransportType() {

2
monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringServiceConfig.java → monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringConfig.java

@ -23,7 +23,7 @@ import org.thingsboard.monitoring.config.TransportType;
@Component
@ConditionalOnProperty(name = "monitoring.transports.http.enabled", havingValue = "true")
@ConfigurationProperties(prefix = "monitoring.transports.http")
public class HttpTransportMonitoringServiceConfig extends TransportMonitoringServiceConfig {
public class HttpTransportMonitoringConfig extends TransportMonitoringConfig {
@Override
public TransportType getTransportType() {

2
monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringServiceConfig.java → monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringConfig.java

@ -27,7 +27,7 @@ import org.thingsboard.monitoring.config.TransportType;
@ConfigurationProperties(prefix = "monitoring.transports.mqtt")
@Data
@EqualsAndHashCode(callSuper = true)
public class MqttTransportMonitoringServiceConfig extends TransportMonitoringServiceConfig {
public class MqttTransportMonitoringConfig extends TransportMonitoringConfig {
private Integer qos;

4
monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringServiceConfig.java → monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringConfig.java

@ -22,11 +22,9 @@ import org.thingsboard.monitoring.config.TransportType;
import java.util.List;
@Data
public abstract class TransportMonitoringServiceConfig {
public abstract class TransportMonitoringConfig {
private int monitoringRateMs;
private int requestTimeoutMs;
private int initialDelayMs;
private List<MonitoringTargetConfig> targets;

2
monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/CmdsWrapper.java

@ -22,6 +22,6 @@ import java.util.List;
@Data
public class CmdsWrapper {
private List<TimeseriesSubscriptionCmd> tsSubCmds;
private List<EntityDataCmd> entityDataCmds;
}

15
monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesSubscriptionCmd.java → monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataCmd.java

@ -15,21 +15,14 @@
*/
package org.thingsboard.monitoring.data.cmd;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.query.EntityDataQuery;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class TimeseriesSubscriptionCmd {
public class EntityDataCmd {
private int cmdId;
private String entityType;
private String entityId;
private String keys;
private String scope;
private EntityDataQuery query;
private LatestValueCmd latestCmd;
}

44
monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java

@ -0,0 +1,44 @@
/**
* Copyright © 2016-2022 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.monitoring.data.cmd;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import lombok.Data;
import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.TsValue;
import java.util.List;
import java.util.UUID;
@Data
@JsonIgnoreProperties(ignoreUnknown = true)
public class EntityDataUpdate {
@JsonIgnoreProperties(ignoreUnknown = true)
private List<EntityData> update;
public String getLatest(UUID entityId, String key) {
if (update == null) return null;
return update.stream()
.filter(entityData -> entityData.getEntityId().getId().equals(entityId)).findFirst()
.map(EntityData::getLatest).map(latest -> latest.get(EntityKeyType.TIME_SERIES))
.map(latest -> latest.get(key)).map(TsValue::getValue)
.orElse(null);
}
}

28
monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/LatestValueCmd.java

@ -0,0 +1,28 @@
/**
* Copyright © 2016-2022 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.monitoring.data.cmd;
import lombok.Data;
import org.thingsboard.server.common.data.query.EntityKey;
import java.util.List;
@Data
public class LatestValueCmd {
private List<EntityKey> keys;
}

52
monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java

@ -1,52 +0,0 @@
/**
* Copyright © 2016-2022 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.monitoring.data.cmd;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import lombok.Data;
import org.springframework.data.util.Pair;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@Data
@JsonIgnoreProperties(ignoreUnknown = true)
public class TimeseriesUpdate {
private final int subscriptionId;
private final int errorCode;
private final String errorMsg;
private Map<String, List<List<Object>>> data;
public Object getLatest(String key) {
if (!data.containsKey(key)) return null;
return data.get(key).stream()
.map(tsAndValue -> {
if (tsAndValue == null || tsAndValue.size() != 2) {
return null;
}
long ts = Long.parseLong(tsAndValue.get(0).toString());
Object value = tsAndValue.get(1);
return Pair.of(ts, value);
})
.filter(Objects::nonNull)
.max(Comparator.comparing(Pair::getFirst))
.map(Pair::getSecond).orElse(null);
}
}

74
monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java

@ -27,6 +27,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.monitoring.client.TbClient;
import org.thingsboard.monitoring.data.Latency;
import org.thingsboard.monitoring.data.MonitoredServiceKey;
import org.thingsboard.monitoring.data.notification.HighLatencyNotification;
import org.thingsboard.monitoring.data.notification.ServiceFailureNotification;
import org.thingsboard.monitoring.data.notification.ServiceRecoveryNotification;
@ -48,9 +49,7 @@ import java.util.stream.Collectors;
@Slf4j
public class MonitoringReporter {
private final TbClient tbClient;
private final NotificationService notificationService;
private final ScheduledExecutorService monitoringExecutor;
private final Map<String, Latency> latencies = new ConcurrentHashMap<>();
private final Map<Object, AtomicInteger> failuresCounters = new ConcurrentHashMap<>();
@ -65,48 +64,43 @@ public class MonitoringReporter {
private EntityType reportingEntityType;
@Value("${monitoring.latency.reporting_entity_id}")
private String reportingEntityId;
@Value("${monitoring.latency.monitoring_rate_ms}")
private int latenciesMonitoringRateMs;
@EventListener(ApplicationReadyEvent.class)
public void startLatenciesMonitoring() {
monitoringExecutor.scheduleAtFixedRate(() -> {
List<Latency> latencies = this.latencies.values().stream()
.filter(Latency::isNotEmpty)
.map(latency -> {
Latency snapshot = latency.snapshot();
latency.reset();
return snapshot;
})
.collect(Collectors.toList());
if (latencies.isEmpty()) {
return;
}
log.info("Latencies:\n{}", latencies);
if (latencies.stream().anyMatch(latency -> latency.getAvg() >= (double) latencyThresholdMs)) {
HighLatencyNotification highLatencyNotification = new HighLatencyNotification(latencies, latencyThresholdMs);
notificationService.sendNotification(highLatencyNotification);
}
public void reportLatencies(TbClient tbClient) {
List<Latency> latencies = this.latencies.values().stream()
.filter(Latency::isNotEmpty)
.map(latency -> {
Latency snapshot = latency.snapshot();
latency.reset();
return snapshot;
})
.collect(Collectors.toList());
if (latencies.isEmpty()) {
return;
}
log.info("Latencies:\n{}", latencies.stream().map(latency -> latency.getKey() + ": " + latency.getAvg() + " ms")
.collect(Collectors.joining("\n")));
if (latencies.stream().anyMatch(latency -> latency.getAvg() >= (double) latencyThresholdMs)) {
HighLatencyNotification highLatencyNotification = new HighLatencyNotification(latencies, latencyThresholdMs);
notificationService.sendNotification(highLatencyNotification);
}
if (reportingEntityType != null && StringUtils.isNotBlank(reportingEntityId)) {
if (reportingEntityType != null && StringUtils.isNotBlank(reportingEntityId)) {
try {
EntityId entityId;
try {
EntityId entityId;
try {
entityId = EntityIdFactory.getByTypeAndUuid(reportingEntityType, reportingEntityId);
} catch (Exception e) {
return;
}
tbClient.logIn();
ObjectNode msg = JacksonUtil.newObjectNode();
latencies.forEach(latency -> {
msg.set(latency.getKey(), new DoubleNode(latency.getAvg()));
});
tbClient.saveEntityTelemetry(entityId, "time", msg);
entityId = EntityIdFactory.getByTypeAndUuid(reportingEntityType, reportingEntityId);
} catch (Exception e) {
log.error("Failed to report latencies: {}", e.getMessage());
return;
}
ObjectNode msg = JacksonUtil.newObjectNode();
latencies.forEach(latency -> {
msg.set(latency.getKey(), new DoubleNode(latency.getAvg()));
});
tbClient.saveEntityTelemetry(entityId, "time", msg);
} catch (Exception e) {
log.error("Failed to report latencies: {}", e.getMessage());
}
}, latenciesMonitoringRateMs, latenciesMonitoringRateMs, TimeUnit.MILLISECONDS);
}
}
public void reportLatency(String key, long latencyInNanos) {
@ -127,7 +121,9 @@ public class MonitoringReporter {
public void serviceIsOk(Object serviceKey) {
ServiceRecoveryNotification notification = new ServiceRecoveryNotification(serviceKey);
log.info(notification.getText());
if (!serviceKey.equals(MonitoredServiceKey.GENERAL)) {
log.info(notification.getText());
}
AtomicInteger failuresCounter = failuresCounters.get(serviceKey);
if (failuresCounter != null) {
if (failuresCounter.get() >= failuresThreshold) {

65
monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java → monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportHealthChecker.java

@ -13,34 +13,30 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.monitoring.service;
package org.thingsboard.monitoring.transport;
import com.fasterxml.jackson.databind.node.TextNode;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.monitoring.client.TbClient;
import org.thingsboard.monitoring.client.WsClient;
import org.thingsboard.monitoring.client.WsClientFactory;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.TransportType;
import org.thingsboard.monitoring.config.WsConfig;
import org.thingsboard.monitoring.config.service.TransportMonitoringServiceConfig;
import org.thingsboard.monitoring.config.service.TransportMonitoringConfig;
import org.thingsboard.monitoring.data.Latencies;
import org.thingsboard.monitoring.data.MonitoredServiceKey;
import org.thingsboard.monitoring.data.TransportFailureException;
import org.thingsboard.monitoring.data.TransportInfo;
import org.thingsboard.monitoring.service.MonitoringReporter;
import org.thingsboard.monitoring.util.TbStopWatch;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@Slf4j
public abstract class TransportMonitoringService<C extends TransportMonitoringServiceConfig> {
public abstract class TransportHealthChecker<C extends TransportMonitoringConfig> {
protected final C config;
protected final MonitoringTargetConfig target;
@ -49,19 +45,13 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
private WsConfig wsConfig;
@Autowired
private MonitoringReporter monitoringReporter;
@Autowired
private WsClientFactory wsClientFactory;
@Autowired
private ScheduledExecutorService monitoringExecutor;
@Autowired
private TbClient tbClient;
private MonitoringReporter reporter;
@Autowired
private TbStopWatch stopWatch;
protected static final String TEST_TELEMETRY_KEY = "testData";
public static final String TEST_TELEMETRY_KEY = "testData";
protected TransportMonitoringService(C config, MonitoringTargetConfig target) {
protected TransportHealthChecker(C config, MonitoringTargetConfig target) {
this.config = config;
this.target = target;
}
@ -71,16 +61,9 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
transportInfo = new TransportInfo(getTransportType(), target.getBaseUrl());
}
public final void startMonitoring() {
monitoringExecutor.scheduleWithFixedDelay(() -> {
check();
}, config.getInitialDelayMs(), config.getMonitoringRateMs(), TimeUnit.MILLISECONDS);
log.info("Started monitoring for transport type {} for target {}", getTransportType(), target);
}
private void check() {
log.trace("[{}] Checking", transportInfo);
try (WsClient wsClient = establishWsClient()) {
public final void check(WsClient wsClient) {
log.debug("[{}] Checking", transportInfo);
try {
wsClient.registerWaitForUpdate();
String testValue = UUID.randomUUID().toString();
@ -95,12 +78,12 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
log.trace("[{}] Waiting for WS update", transportInfo);
checkWsUpdate(wsClient, testValue);
monitoringReporter.serviceIsOk(transportInfo);
monitoringReporter.serviceIsOk(MonitoredServiceKey.GENERAL);
reporter.serviceIsOk(transportInfo);
reporter.serviceIsOk(MonitoredServiceKey.GENERAL);
} catch (TransportFailureException transportFailureException) {
monitoringReporter.serviceFailure(transportInfo, transportFailureException);
reporter.serviceFailure(transportInfo, transportFailureException);
} catch (Exception e) {
monitoringReporter.serviceFailure(MonitoredServiceKey.GENERAL, e);
reporter.serviceFailure(MonitoredServiceKey.GENERAL, e);
}
}
@ -108,34 +91,20 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
initClient();
stopWatch.start();
sendTestPayload(payload);
monitoringReporter.reportLatency(Latencies.transportRequest(getTransportType()), stopWatch.getTime());
}
private WsClient establishWsClient() throws Exception {
stopWatch.start();
String accessToken = tbClient.logIn();
log.trace("[{}] Received new access token", transportInfo);
monitoringReporter.reportLatency(Latencies.LOG_IN, stopWatch.getTime());
WsClient wsClient = wsClientFactory.createClient(accessToken);
log.trace("[{}] Created WS client", transportInfo);
wsClient.subscribeForTelemetry(target.getDevice().getId(), TEST_TELEMETRY_KEY);
Optional.ofNullable(wsClient.waitForReply(wsConfig.getRequestTimeoutMs()))
.orElseThrow(() -> new IllegalStateException("Failed to subscribe for telemetry"));
return wsClient;
reporter.reportLatency(Latencies.transportRequest(getTransportType()), stopWatch.getTime());
}
private void checkWsUpdate(WsClient wsClient, String testValue) {
stopWatch.start();
wsClient.waitForUpdate(wsConfig.getResultCheckTimeoutMs());
log.trace("[{}] Waited for WS update. Last WS msg: {}", transportInfo, wsClient.lastMsg);
Object update = wsClient.getTelemetryKeyUpdate(TEST_TELEMETRY_KEY);
Object update = wsClient.getTelemetryUpdate(target.getDevice().getId(), TEST_TELEMETRY_KEY);
if (update == null) {
throw new TransportFailureException("No WS update arrived within " + wsConfig.getResultCheckTimeoutMs() + " ms");
} else if (!update.toString().equals(testValue)) {
throw new TransportFailureException("Was expecting value " + testValue + " but got " + update);
}
monitoringReporter.reportLatency(Latencies.WS_UPDATE, stopWatch.getTime());
reporter.reportLatency(Latencies.WS_UPDATE, stopWatch.getTime());
}

143
monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportMonitoringService.java

@ -0,0 +1,143 @@
/**
* Copyright © 2016-2022 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.monitoring.transport;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.ApplicationContext;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.monitoring.client.TbClient;
import org.thingsboard.monitoring.client.WsClient;
import org.thingsboard.monitoring.client.WsClientFactory;
import org.thingsboard.monitoring.config.DeviceConfig;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.service.TransportMonitoringConfig;
import org.thingsboard.monitoring.data.Latencies;
import org.thingsboard.monitoring.data.MonitoredServiceKey;
import org.thingsboard.monitoring.service.MonitoringReporter;
import org.thingsboard.monitoring.util.TbStopWatch;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration;
import org.thingsboard.server.common.data.device.data.DefaultDeviceTransportConfiguration;
import org.thingsboard.server.common.data.device.data.DeviceData;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import javax.annotation.PostConstruct;
import java.util.LinkedList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@Service
@RequiredArgsConstructor
@Slf4j
public final class TransportMonitoringService {
private final List<TransportMonitoringConfig> configs;
private final List<TransportHealthChecker<?>> transportHealthCheckers = new LinkedList<>();
private final List<UUID> devices = new LinkedList<>();
private final TbClient tbClient;
private final WsClientFactory wsClientFactory;
private final TbStopWatch stopWatch;
private final MonitoringReporter reporter;
private final ApplicationContext applicationContext;
private ScheduledExecutorService scheduler;
@Value("${monitoring.transports.monitoring_rate_ms}")
private int monitoringRateMs;
@PostConstruct
private void init() {
configs.forEach(config -> {
config.getTargets().stream()
.filter(target -> StringUtils.isNotBlank(target.getBaseUrl()))
.peek(target -> checkMonitoringTarget(config, target, tbClient))
.forEach(target -> {
TransportHealthChecker<?> transportHealthChecker = applicationContext.getBean(config.getTransportType().getServiceClass(), config, target);
transportHealthCheckers.add(transportHealthChecker);
devices.add(target.getDevice().getId());
});
});
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("monitoring-executor"));
}
@EventListener(ApplicationReadyEvent.class)
public void startMonitoring() {
scheduleCheck(0);
}
private void scheduleCheck(int delay) {
log.debug("Scheduling next check for {} ms", delay);
scheduler.schedule(() -> {
try {
stopWatch.start();
String accessToken = tbClient.logIn();
reporter.reportLatency(Latencies.LOG_IN, stopWatch.getTime());
try (WsClient wsClient = wsClientFactory.createClient(accessToken)) {
wsClient.subscribeForTelemetry(devices, TransportHealthChecker.TEST_TELEMETRY_KEY).waitForReply();
for (TransportHealthChecker<?> transportHealthChecker : transportHealthCheckers) {
transportHealthChecker.check(wsClient);
}
}
reporter.reportLatencies(tbClient);
} catch (Exception e) {
reporter.serviceFailure(MonitoredServiceKey.GENERAL, e);
}
scheduleCheck(monitoringRateMs);
}, delay, TimeUnit.MILLISECONDS);
}
private void checkMonitoringTarget(TransportMonitoringConfig config, MonitoringTargetConfig target, TbClient tbClient) {
DeviceConfig deviceConfig = target.getDevice();
tbClient.logIn();
DeviceId deviceId;
if (deviceConfig == null || deviceConfig.getId() == null) {
String deviceName = String.format("[%s] Monitoring device (%s)", config.getTransportType(), target.getBaseUrl());
Device device = tbClient.getTenantDevice(deviceName)
.orElseGet(() -> {
log.info("Creating new device '{}'", deviceName);
Device monitoringDevice = new Device();
monitoringDevice.setName(deviceName);
monitoringDevice.setType("default");
DeviceData deviceData = new DeviceData();
deviceData.setConfiguration(new DefaultDeviceConfiguration());
deviceData.setTransportConfiguration(new DefaultDeviceTransportConfiguration());
return tbClient.saveDevice(monitoringDevice);
});
deviceId = device.getId();
target.getDevice().setId(deviceId.toString());
} else {
deviceId = new DeviceId(deviceConfig.getId());
}
log.debug("Loading credentials for device {}", deviceId);
DeviceCredentials credentials = tbClient.getDeviceCredentialsByDeviceId(deviceId)
.orElseThrow(() -> new IllegalArgumentException("No credentials found for device " + deviceId));
target.getDevice().setCredentials(credentials);
}
}

10
monitoring/src/main/java/org/thingsboard/monitoring/service/impl/CoapTransportMonitoringService.java → monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/CoapTransportHealthChecker.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.monitoring.service.impl;
package org.thingsboard.monitoring.transport.impl;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.core.CoapClient;
@ -25,19 +25,19 @@ import org.springframework.context.annotation.Scope;
import org.springframework.stereotype.Component;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.TransportType;
import org.thingsboard.monitoring.config.service.CoapTransportMonitoringServiceConfig;
import org.thingsboard.monitoring.service.TransportMonitoringService;
import org.thingsboard.monitoring.config.service.CoapTransportMonitoringConfig;
import org.thingsboard.monitoring.transport.TransportHealthChecker;
import java.io.IOException;
@Component
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Slf4j
public class CoapTransportMonitoringService extends TransportMonitoringService<CoapTransportMonitoringServiceConfig> {
public class CoapTransportHealthChecker extends TransportHealthChecker<CoapTransportMonitoringConfig> {
private CoapClient coapClient;
protected CoapTransportMonitoringService(CoapTransportMonitoringServiceConfig config, MonitoringTargetConfig target) {
protected CoapTransportHealthChecker(CoapTransportMonitoringConfig config, MonitoringTargetConfig target) {
super(config, target);
}

10
monitoring/src/main/java/org/thingsboard/monitoring/service/impl/HttpTransportMonitoringService.java → monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/HttpTransportHealthChecker.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.monitoring.service.impl;
package org.thingsboard.monitoring.transport.impl;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.config.ConfigurableBeanFactory;
@ -23,19 +23,19 @@ import org.springframework.stereotype.Component;
import org.springframework.web.client.RestTemplate;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.TransportType;
import org.thingsboard.monitoring.config.service.HttpTransportMonitoringServiceConfig;
import org.thingsboard.monitoring.service.TransportMonitoringService;
import org.thingsboard.monitoring.config.service.HttpTransportMonitoringConfig;
import org.thingsboard.monitoring.transport.TransportHealthChecker;
import java.time.Duration;
@Component
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Slf4j
public class HttpTransportMonitoringService extends TransportMonitoringService<HttpTransportMonitoringServiceConfig> {
public class HttpTransportHealthChecker extends TransportHealthChecker<HttpTransportMonitoringConfig> {
private RestTemplate restTemplate;
protected HttpTransportMonitoringService(HttpTransportMonitoringServiceConfig config, MonitoringTargetConfig target) {
protected HttpTransportHealthChecker(HttpTransportMonitoringConfig config, MonitoringTargetConfig target) {
super(config, target);
}

10
monitoring/src/main/java/org/thingsboard/monitoring/service/impl/MqttTransportMonitoringService.java → monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/MqttTransportHealthChecker.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.monitoring.service.impl;
package org.thingsboard.monitoring.transport.impl;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.IMqttToken;
@ -27,19 +27,19 @@ import org.springframework.context.annotation.Scope;
import org.springframework.stereotype.Component;
import org.thingsboard.monitoring.config.MonitoringTargetConfig;
import org.thingsboard.monitoring.config.TransportType;
import org.thingsboard.monitoring.config.service.MqttTransportMonitoringServiceConfig;
import org.thingsboard.monitoring.service.TransportMonitoringService;
import org.thingsboard.monitoring.config.service.MqttTransportMonitoringConfig;
import org.thingsboard.monitoring.transport.TransportHealthChecker;
@Component
@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)
@Slf4j
public class MqttTransportMonitoringService extends TransportMonitoringService<MqttTransportMonitoringServiceConfig> {
public class MqttTransportHealthChecker extends TransportHealthChecker<MqttTransportMonitoringConfig> {
private MqttClient mqttClient;
private static final String DEVICE_TELEMETRY_TOPIC = "v1/devices/me/telemetry";
protected MqttTransportMonitoringService(MqttTransportMonitoringServiceConfig config, MonitoringTargetConfig target) {
protected MqttTransportHealthChecker(MqttTransportMonitoringConfig config, MonitoringTargetConfig target) {
super(config, target);
}

2
monitoring/src/main/resources/logback.xml

@ -29,8 +29,6 @@
<logger name="org.thingsboard.server" level="INFO"/>
<logger name="org.thingsboard.monitoring" level="DEBUG"/>
<logger name="org.thingsboard.monitoring.client" level="WARN"/>
<logger name="org.openqa.selenium" level="WARN"/>
<logger name="io.github.bonigarcia" level="WARN"/>
<root level="INFO">
<appender-ref ref="STDOUT"/>

11
monitoring/src/main/resources/tb-monitoring.yml

@ -32,11 +32,11 @@ monitoring:
send_repeated_failure_notification: '${SEND_REPEATED_FAILURE_NOTIFICATION:true}'
transports:
monitoring_rate_ms: '${TRANSPORTS_MONITORING_RATE_MS:10000}'
mqtt:
enabled: '${MQTT_TRANSPORT_MONITORING_ENABLED:false}'
monitoring_rate_ms: '${MQTT_TRANSPORT_MONITORING_RATE_MS:3000}'
request_timeout_ms: '${MQTT_REQUEST_TIMEOUT_MS:4000}'
initial_delay_ms: '${MQTT_TRANSPORT_MONITORING_INITIAL_DELAY_MS:0}'
qos: '${MQTT_QOS_LEVEL:1}'
targets:
- base_url: '${MQTT_TRANSPORT_BASE_URL:tcp://localhost:1883}'
@ -45,9 +45,7 @@ monitoring:
coap:
enabled: '${COAP_TRANSPORT_MONITORING_ENABLED:false}'
monitoring_rate_ms: '${COAP_TRANSPORT_MONITORING_RATE_MS:10000}'
request_timeout_ms: '${COAP_REQUEST_TIMEOUT_MS:4000}'
initial_delay_ms: '${COAP_TRANSPORT_MONITORING_INITIAL_DELAY_MS:0}'
targets:
- base_url: '${COAP_TRANSPORT_BASE_URL:coap://localhost}'
device:
@ -55,9 +53,7 @@ monitoring:
http:
enabled: '${HTTP_TRANSPORT_MONITORING_ENABLED:false}'
monitoring_rate_ms: '${HTTP_TRANSPORT_MONITORING_RATE_MS:10000}'
request_timeout_ms: '${HTTP_REQUEST_TIMEOUT_MS:4000}'
initial_delay_ms: '${HTTP_TRANSPORT_MONITORING_INITIAL_DELAY_MS:0}'
targets:
- base_url: '${HTTP_TRANSPORT_BASE_URL:http://localhost:8080}'
device:
@ -69,9 +65,6 @@ monitoring:
webhook_url: '${SLACK_WEBHOOK_URL:}'
latency:
monitoring_rate_ms: '${LATENCY_MONITORING_RATE_MS:30000}'
threshold_ms: '${LATENCY_THRESHOLD:2000}'
reporting_entity_type: '${LATENCY_REPORTING_ENTITY_TYPE:ASSET}'
reporting_entity_id: '${LATENCY_REPORTING_ENTITY_ID:}'
monitoring_executor_thread_pool_size: '${MONITORING_EXECUTOR_THREAD_POOL_SIZE:10}'

Loading…
Cancel
Save