From 869136da8909dfb259cd17cb9cb445c990c7f00b Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 28 Feb 2023 13:38:50 +0200 Subject: [PATCH] Monitoring service to single-thread --- .../ThingsboardMonitoringApplication.java | 70 --------- .../monitoring/client/WsClient.java | 72 +++++---- .../monitoring/client/WsClientFactory.java | 2 +- .../monitoring/config/TransportType.java | 17 ++- ...ava => CoapTransportMonitoringConfig.java} | 2 +- ...ava => HttpTransportMonitoringConfig.java} | 2 +- ...ava => MqttTransportMonitoringConfig.java} | 2 +- ...ig.java => TransportMonitoringConfig.java} | 4 +- .../monitoring/data/cmd/CmdsWrapper.java | 2 +- ...ubscriptionCmd.java => EntityDataCmd.java} | 15 +- .../monitoring/data/cmd/EntityDataUpdate.java | 44 ++++++ .../monitoring/data/cmd/LatestValueCmd.java | 28 ++++ .../monitoring/data/cmd/TimeseriesUpdate.java | 52 ------- .../service/MonitoringReporter.java | 74 +++++---- .../TransportHealthChecker.java} | 65 +++----- .../transport/TransportMonitoringService.java | 143 ++++++++++++++++++ .../impl/CoapTransportHealthChecker.java} | 10 +- .../impl/HttpTransportHealthChecker.java} | 10 +- .../impl/MqttTransportHealthChecker.java} | 10 +- monitoring/src/main/resources/logback.xml | 2 - .../src/main/resources/tb-monitoring.yml | 11 +- 21 files changed, 348 insertions(+), 289 deletions(-) rename monitoring/src/main/java/org/thingsboard/monitoring/config/service/{CoapTransportMonitoringServiceConfig.java => CoapTransportMonitoringConfig.java} (92%) rename monitoring/src/main/java/org/thingsboard/monitoring/config/service/{HttpTransportMonitoringServiceConfig.java => HttpTransportMonitoringConfig.java} (92%) rename monitoring/src/main/java/org/thingsboard/monitoring/config/service/{MqttTransportMonitoringServiceConfig.java => MqttTransportMonitoringConfig.java} (93%) rename monitoring/src/main/java/org/thingsboard/monitoring/config/service/{TransportMonitoringServiceConfig.java => TransportMonitoringConfig.java} (88%) rename monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/{TimeseriesSubscriptionCmd.java => EntityDataCmd.java} (71%) create mode 100644 monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java create mode 100644 monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/LatestValueCmd.java delete mode 100644 monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java rename monitoring/src/main/java/org/thingsboard/monitoring/{service/TransportMonitoringService.java => transport/TransportHealthChecker.java} (59%) create mode 100644 monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportMonitoringService.java rename monitoring/src/main/java/org/thingsboard/monitoring/{service/impl/CoapTransportMonitoringService.java => transport/impl/CoapTransportHealthChecker.java} (87%) rename monitoring/src/main/java/org/thingsboard/monitoring/{service/impl/HttpTransportMonitoringService.java => transport/impl/HttpTransportHealthChecker.java} (85%) rename monitoring/src/main/java/org/thingsboard/monitoring/{service/impl/MqttTransportMonitoringService.java => transport/impl/MqttTransportHealthChecker.java} (89%) diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java b/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java index aa9d24e319..8997c0d373 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java +++ b/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 configs, TbClient tbClient, ApplicationContext context) { - return args -> { - List> 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")); - } - } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java b/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java index ea1a5035f2..7d2cbcd936 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java +++ b/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 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 diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClientFactory.java b/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClientFactory.java index 1ed9fb8a24..b3fc4f1558 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClientFactory.java +++ b/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); diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/TransportType.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/TransportType.java index 49a9cda010..8b2ec3151b 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/TransportType.java +++ b/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> monitoringServiceClass; + MQTT(MqttTransportHealthChecker.class), + COAP(CoapTransportHealthChecker.class), + HTTP(HttpTransportHealthChecker.class); + + private final Class> serviceClass; } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringServiceConfig.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringConfig.java similarity index 92% rename from monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringServiceConfig.java rename to monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringConfig.java index 9218df12ee..ff6b7b25bd 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/CoapTransportMonitoringServiceConfig.java +++ b/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() { diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringServiceConfig.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringConfig.java similarity index 92% rename from monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringServiceConfig.java rename to monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringConfig.java index 90f7d86957..a7250e4f49 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/HttpTransportMonitoringServiceConfig.java +++ b/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() { diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringServiceConfig.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringConfig.java similarity index 93% rename from monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringServiceConfig.java rename to monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringConfig.java index b2e58510ff..e8c34f7db7 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/MqttTransportMonitoringServiceConfig.java +++ b/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; diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringServiceConfig.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringConfig.java similarity index 88% rename from monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringServiceConfig.java rename to monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringConfig.java index 16c4bd1473..d2d75a888a 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/service/TransportMonitoringServiceConfig.java +++ b/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 targets; diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/CmdsWrapper.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/CmdsWrapper.java index cc0c482d3e..3d39f711a2 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/CmdsWrapper.java +++ b/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 tsSubCmds; + private List entityDataCmds; } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesSubscriptionCmd.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataCmd.java similarity index 71% rename from monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesSubscriptionCmd.java rename to monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataCmd.java index 2f3d65647a..b0990b70f6 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesSubscriptionCmd.java +++ b/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; } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java new file mode 100644 index 0000000000..debb123e15 --- /dev/null +++ b/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 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); + } + +} diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/LatestValueCmd.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/LatestValueCmd.java new file mode 100644 index 0000000000..a09c41c184 --- /dev/null +++ b/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 keys; + +} diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java deleted file mode 100644 index 515fae2c1d..0000000000 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/TimeseriesUpdate.java +++ /dev/null @@ -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>> 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); - } - -} diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java b/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java index 5721609d27..3b954efd2f 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java +++ b/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 latencies = new ConcurrentHashMap<>(); private final Map 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 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 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) { diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportHealthChecker.java similarity index 59% rename from monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java rename to monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportHealthChecker.java index 4fe003d87d..f73f935664 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java +++ b/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 { +public abstract class TransportHealthChecker { protected final C config; protected final MonitoringTargetConfig target; @@ -49,19 +45,13 @@ public abstract class TransportMonitoringService { - 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 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()); } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/transport/TransportMonitoringService.java new file mode 100644 index 0000000000..2d6e0597b6 --- /dev/null +++ b/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 configs; + private final List> transportHealthCheckers = new LinkedList<>(); + private final List 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); + } + +} diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/CoapTransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/CoapTransportHealthChecker.java similarity index 87% rename from monitoring/src/main/java/org/thingsboard/monitoring/service/impl/CoapTransportMonitoringService.java rename to monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/CoapTransportHealthChecker.java index 43558a6d49..84568623c4 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/CoapTransportMonitoringService.java +++ b/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 { +public class CoapTransportHealthChecker extends TransportHealthChecker { private CoapClient coapClient; - protected CoapTransportMonitoringService(CoapTransportMonitoringServiceConfig config, MonitoringTargetConfig target) { + protected CoapTransportHealthChecker(CoapTransportMonitoringConfig config, MonitoringTargetConfig target) { super(config, target); } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/HttpTransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/HttpTransportHealthChecker.java similarity index 85% rename from monitoring/src/main/java/org/thingsboard/monitoring/service/impl/HttpTransportMonitoringService.java rename to monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/HttpTransportHealthChecker.java index 1062d675d6..f055e9958d 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/HttpTransportMonitoringService.java +++ b/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 { +public class HttpTransportHealthChecker extends TransportHealthChecker { private RestTemplate restTemplate; - protected HttpTransportMonitoringService(HttpTransportMonitoringServiceConfig config, MonitoringTargetConfig target) { + protected HttpTransportHealthChecker(HttpTransportMonitoringConfig config, MonitoringTargetConfig target) { super(config, target); } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/MqttTransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/MqttTransportHealthChecker.java similarity index 89% rename from monitoring/src/main/java/org/thingsboard/monitoring/service/impl/MqttTransportMonitoringService.java rename to monitoring/src/main/java/org/thingsboard/monitoring/transport/impl/MqttTransportHealthChecker.java index 5cab9f146f..ca8ba275aa 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/impl/MqttTransportMonitoringService.java +++ b/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 { +public class MqttTransportHealthChecker extends TransportHealthChecker { 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); } diff --git a/monitoring/src/main/resources/logback.xml b/monitoring/src/main/resources/logback.xml index 945c7dd18e..99e2c7172d 100644 --- a/monitoring/src/main/resources/logback.xml +++ b/monitoring/src/main/resources/logback.xml @@ -29,8 +29,6 @@ - - diff --git a/monitoring/src/main/resources/tb-monitoring.yml b/monitoring/src/main/resources/tb-monitoring.yml index 3d71a7dbfa..d0f9be4ea2 100644 --- a/monitoring/src/main/resources/tb-monitoring.yml +++ b/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}'