From 82d4cb538112a98ed6425cf56c63e27cceb6731f Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 15 Apr 2025 15:12:18 +0300 Subject: [PATCH] Add monitoring for calculated fields --- .../monitoring/client/WsClient.java | 69 +++++++++++-------- .../config/transport/TransportInfo.java | 11 ++- .../transport/TransportMonitoringTarget.java | 1 + .../monitoring/data/MonitoredServiceKey.java | 1 + .../data/ServiceFailureException.java | 11 ++- .../monitoring/data/cmd/EntityDataUpdate.java | 18 +++-- .../monitoring/service/BaseHealthChecker.java | 44 ++++++++---- .../service/BaseMonitoringService.java | 37 +++++----- .../service/MonitoringReporter.java | 2 +- .../transport/TransportHealthChecker.java | 44 +++++++++++- .../src/main/resources/tb-monitoring.yml | 11 +++ .../thingsboard/rest/client/RestClient.java | 20 +++++- 12 files changed, 189 insertions(+), 80 deletions(-) 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 9246feac27..9bacbdfd45 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java @@ -36,8 +36,11 @@ 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.ArrayList; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -48,13 +51,13 @@ import java.util.stream.Collectors; @Slf4j public class WsClient extends WebSocketClient implements AutoCloseable { - public volatile JsonNode lastMsg; + public final List lastMsgs = new ArrayList<>(); private CountDownLatch reply; private CountDownLatch update; private final Lock updateLock = new ReentrantLock(); - private long requestTimeoutMs; + private final long requestTimeoutMs; public WsClient(URI serverUri, long requestTimeoutMs) { super(serverUri); @@ -63,7 +66,6 @@ public class WsClient extends WebSocketClient implements AutoCloseable { @Override public void onOpen(ServerHandshake serverHandshake) { - } @Override @@ -73,8 +75,9 @@ public class WsClient extends WebSocketClient implements AutoCloseable { } updateLock.lock(); try { - lastMsg = JacksonUtil.toJsonNode(s); - log.trace("Received new msg: {}", lastMsg.toPrettyString()); + JsonNode msg = JacksonUtil.toJsonNode(s); + lastMsgs.add(msg); + log.trace("Received new msg: {}", msg.toPrettyString()); if (update != null) { update.countDown(); } @@ -96,11 +99,11 @@ public class WsClient extends WebSocketClient implements AutoCloseable { log.error("WebSocket client error:", e); } - public void registerWaitForUpdate() { + public void registerWaitForUpdates(int count) { updateLock.lock(); try { - lastMsg = null; - update = new CountDownLatch(1); + lastMsgs.clear(); + update = new CountDownLatch(count); } finally { updateLock.unlock(); } @@ -111,6 +114,7 @@ public class WsClient extends WebSocketClient implements AutoCloseable { public void send(String text) throws NotYetConnectedException { updateLock.lock(); try { + lastMsgs.clear(); reply = new CountDownLatch(1); } finally { updateLock.unlock(); @@ -118,19 +122,19 @@ public class WsClient extends WebSocketClient implements AutoCloseable { super.send(text); } - public WsClient subscribeForTelemetry(List devices, String key) { + public WsClient subscribeForTelemetry(List devices, List keys) { 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); + 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))); + latestCmd.setKeys(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).toList()); cmd.setLatestCmd(latestCmd); CmdsWrapper wrapper = new CmdsWrapper(); @@ -139,12 +143,12 @@ public class WsClient extends WebSocketClient implements AutoCloseable { return this; } - public JsonNode waitForUpdate(long ms) { + public List waitForUpdates(long ms) { log.trace("update latch count: {}", update.getCount()); try { if (update.await(ms, TimeUnit.MILLISECONDS)) { log.trace("Waited for update"); - return getLastMsg(); + return getLastMsgs(); } } catch (InterruptedException e) { log.debug("Failed to await reply", e); @@ -157,7 +161,8 @@ public class WsClient extends WebSocketClient implements AutoCloseable { try { if (reply.await(requestTimeoutMs, TimeUnit.MILLISECONDS)) { log.trace("Waited for reply"); - return getLastMsg(); + List lastMsgs = getLastMsgs(); + return lastMsgs.isEmpty() ? null : lastMsgs.get(0); } } catch (InterruptedException e) { log.debug("Failed to await reply", e); @@ -166,24 +171,30 @@ public class WsClient extends WebSocketClient implements AutoCloseable { throw new IllegalStateException("No WS reply arrived within " + requestTimeoutMs + " ms"); } - 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 lastMsg; - } - } else { - return null; + private List getLastMsgs() { + if (lastMsgs.isEmpty()) { + return lastMsgs; + } + List errors = lastMsgs.stream() + .map(msg -> msg.get("errorMsg")) + .filter(errorMsg -> errorMsg != null && !errorMsg.isNull() && StringUtils.isNotEmpty(errorMsg.asText())) + .toList(); + if (!errors.isEmpty()) { + throw new RuntimeException("WS error from server: " + errors.stream() + .map(JsonNode::asText) + .collect(Collectors.joining(", "))); } + return lastMsgs; } - public Object getTelemetryUpdate(UUID deviceId, String key) { - JsonNode lastMsg = getLastMsg(); - if (lastMsg == null || lastMsg.isNull()) return null; - EntityDataUpdate update = JacksonUtil.treeToValue(lastMsg, EntityDataUpdate.class); - return update.getLatest(deviceId, key); + public Map getLatest(UUID deviceId) { + Map updates = new HashMap<>(); + getLastMsgs().forEach(msg -> { + EntityDataUpdate update = JacksonUtil.treeToValue(msg, EntityDataUpdate.class); + Map latest = update.getLatest(deviceId); + updates.putAll(latest); + }); + return updates; } @Override diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportInfo.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportInfo.java index 208d28ffee..6b77e1e268 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportInfo.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportInfo.java @@ -20,16 +20,15 @@ import lombok.Data; @Data public class TransportInfo { - private final TransportType transportType; - private final String baseUrl; - private final String queue; + private final TransportType type; + private final TransportMonitoringTarget target; @Override public String toString() { - if (queue.equals("Main")) { - return String.format("*%s* (%s)", transportType.getName(), baseUrl); + if (target.getQueue().equals("Main")) { + return String.format("*%s* (%s)", type.getName(), target.getBaseUrl()); } else { - return String.format("*%s* (%s) _%s_", transportType.getName(), baseUrl, queue); + return String.format("*%s* (%s) _%s_", type.getName(), target.getBaseUrl(), target.getQueue()); } } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java b/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java index 1e6a3c8509..e8a9ab03fa 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java @@ -28,6 +28,7 @@ public class TransportMonitoringTarget implements MonitoringTarget { private DeviceConfig device; // set manually during initialization private String queue; private boolean checkDomainIps; + private String namePrefix; @Override public UUID getDeviceId() { diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/MonitoredServiceKey.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/MonitoredServiceKey.java index 9c3ee5b786..7579d75231 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/MonitoredServiceKey.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/data/MonitoredServiceKey.java @@ -19,5 +19,6 @@ public class MonitoredServiceKey { public static final String GENERAL = "Monitoring"; public static final String EDQS = "*EDQS*"; + public static final String CF = "*CF*"; } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/ServiceFailureException.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/ServiceFailureException.java index 5f46514bd1..b2592b8719 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/ServiceFailureException.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/data/ServiceFailureException.java @@ -15,14 +15,21 @@ */ package org.thingsboard.monitoring.data; +import lombok.Getter; + +@Getter public class ServiceFailureException extends RuntimeException { - public ServiceFailureException(Throwable cause) { + private final Object serviceKey; + + public ServiceFailureException(Object serviceKey, Throwable cause) { super(cause.getMessage(), cause); + this.serviceKey = serviceKey; } - public ServiceFailureException(String message) { + public ServiceFailureException(Object serviceKey, String message) { super(message); + this.serviceKey = serviceKey; } } 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 index b82dcbf54d..0706dedcf3 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/data/cmd/EntityDataUpdate.java @@ -19,9 +19,11 @@ 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.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; @Data @@ -31,14 +33,16 @@ public class EntityDataUpdate { @JsonIgnoreProperties(ignoreUnknown = true) private List update; - public String getLatest(UUID entityId, String key) { - if (update == null) return null; - - return update.stream() + public Map getLatest(UUID entityId) { + if (update == null || update.isEmpty()) { + return Collections.emptyMap(); + } + Map result = new HashMap<>(); + 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); + .ifPresent(latest -> latest.forEach((key, tsValue) -> result.put(key, tsValue.getValue()))); + return result; } } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java b/monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java index 482b3a30fb..34e0ebc4c5 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java @@ -52,11 +52,14 @@ public abstract class BaseHealthChecker> associates = new HashMap<>(); public static final String TEST_TELEMETRY_KEY = "testData"; + public static final String TEST_CF_TELEMETRY_KEY = "testDataCf"; @PostConstruct private void init() { @@ -68,7 +71,8 @@ public abstract class BaseHealthChecker latest = wsClient.getLatest(target.getDeviceId()); + if (latest.isEmpty()) { + throw new ServiceFailureException(info, "No WS update arrived within " + resultCheckTimeoutMs + " ms"); + } + String actualValue = latest.get(TEST_TELEMETRY_KEY); + if (!testValue.equals(actualValue)) { + throw new ServiceFailureException(info, "Was expecting value " + testValue + " but got " + actualValue); + } + if (checkCalculatedFields) { + String cfTestValue = testValue + "-cf"; + String actualCfValue = latest.get(TEST_CF_TELEMETRY_KEY); + if (actualCfValue == null) { + throw new ServiceFailureException(MonitoredServiceKey.CF, "No CF value arrived"); + } else if (!cfTestValue.equals(actualCfValue)) { + throw new ServiceFailureException(MonitoredServiceKey.CF, "Was expecting CF value " + cfTestValue + " but got " + actualCfValue); + } else { + reporter.serviceIsOk(MonitoredServiceKey.CF); + } } reporter.reportLatency(Latencies.wsUpdate(getKey()), stopWatch.getTime()); } @@ -121,6 +138,7 @@ public abstract class BaseHealthChecker, T extends MonitoringTarget> { @@ -79,7 +81,9 @@ public abstract class BaseMonitoringService, T ext protected ApplicationContext applicationContext; @Value("${monitoring.edqs.enabled:false}") - private boolean edqsMonitoringEnabled; + private boolean checkEdqs; + @Value("${monitoring.calculated_fields.enabled:true}") + protected boolean checkCalculatedFields; @PostConstruct private void init() { @@ -121,7 +125,7 @@ public abstract class BaseMonitoringService, T ext try (WsClient wsClient = wsClientFactory.createClient(accessToken)) { stopWatch.start(); - wsClient.subscribeForTelemetry(devices, TransportHealthChecker.TEST_TELEMETRY_KEY).waitForReply(); + wsClient.subscribeForTelemetry(devices, getTestTelemetryKeys()).waitForReply(); reporter.reportLatency(Latencies.WS_SUBSCRIBE, stopWatch.getTime()); for (BaseHealthChecker healthChecker : healthCheckers) { @@ -129,22 +133,17 @@ public abstract class BaseMonitoringService, T ext } } - if (edqsMonitoringEnabled) { - try { - stopWatch.start(); - checkEdqs(); - reporter.reportLatency(Latencies.EDQS_QUERY, stopWatch.getTime()); - - reporter.serviceIsOk(MonitoredServiceKey.EDQS); - } catch (ServiceFailureException e) { - reporter.serviceFailure(MonitoredServiceKey.EDQS, e); - } catch (Exception e) { - reporter.serviceFailure(MonitoredServiceKey.GENERAL, e); - } + if (checkEdqs) { + stopWatch.start(); + checkEdqs(); + reporter.reportLatency(Latencies.EDQS_QUERY, stopWatch.getTime()); + reporter.serviceIsOk(MonitoredServiceKey.EDQS); } reporter.reportLatencies(tbClient); log.debug("Finished {}", getName()); + } catch (ServiceFailureException e) { + reporter.serviceFailure(e.getServiceKey(), e); } catch (Throwable error) { try { reporter.serviceFailure(MonitoredServiceKey.GENERAL, error); @@ -199,7 +198,7 @@ public abstract class BaseMonitoringService, T ext .collect(Collectors.toSet()); Set missing = Sets.difference(new HashSet<>(this.devices), devices); if (!missing.isEmpty()) { - throw new ServiceFailureException("Missing devices in the response: " + missing); + throw new ServiceFailureException(MonitoredServiceKey.EDQS, "Missing devices in the response: " + missing); } result.getData().stream() @@ -211,7 +210,7 @@ public abstract class BaseMonitoringService, T ext Stream.of("name", "type", "testData").forEach(key -> { TsValue value = values.get(key); if (value == null || StringUtils.isBlank(value.getValue())) { - throw new ServiceFailureException("Missing " + key + " for device " + entityData.getEntityId()); + throw new ServiceFailureException(MonitoredServiceKey.EDQS, "Missing " + key + " for device " + entityData.getEntityId()); } }); }); @@ -232,6 +231,10 @@ public abstract class BaseMonitoringService, T ext .collect(Collectors.toSet()); } + private List getTestTelemetryKeys() { + return checkCalculatedFields ? List.of(TEST_TELEMETRY_KEY, TEST_CF_TELEMETRY_KEY) : List.of(TEST_TELEMETRY_KEY); + } + private void stopHealthChecker(BaseHealthChecker healthChecker) throws Exception { healthChecker.destroyClient(); devices.remove(healthChecker.getTarget().getDeviceId()); 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 3554731e58..62ed0d74aa 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java @@ -113,7 +113,7 @@ public class MonitoringReporter { public void serviceFailure(Object serviceKey, Throwable error) { if (log.isDebugEnabled()) { - log.error("Error occurred", error); + log.error("[{}] Error occurred", serviceKey, error); } int failuresCount = failuresCounters.computeIfAbsent(serviceKey, k -> new AtomicInteger()).incrementAndGet(); ServiceFailureNotification notification = new ServiceFailureNotification(serviceKey, error, failuresCount); diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportHealthChecker.java b/monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportHealthChecker.java index f5e6cb5469..b892f5d609 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportHealthChecker.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportHealthChecker.java @@ -32,6 +32,14 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MBootstrapClientCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.NoSecBootstrapClientCredential; @@ -47,6 +55,8 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentialsType; +import java.util.Map; + @Slf4j public abstract class TransportHealthChecker extends BaseHealthChecker { @@ -74,7 +84,7 @@ public abstract class TransportHealthChecker getImages(PageLink pageLink, boolean includeSystemImages) { - return this.getImages(pageLink, null, includeSystemImages); + return this.getImages(pageLink, null, includeSystemImages); } public PageData getImages(PageLink pageLink, ResourceSubType imageSubType, boolean includeSystemImages) { @@ -4056,6 +4057,21 @@ public class RestClient implements Closeable { timeout).getBody(); } + public CalculatedField saveCalculatedField(CalculatedField calculatedField) { + return restTemplate.postForEntity(baseURL + "/api/calculatedField", calculatedField, CalculatedField.class).getBody(); + } + + public PageData getCalculatedFieldsByEntityId(EntityId entityId, PageLink pageLink) { + Map params = new HashMap<>(); + addPageLinkToParam(params, pageLink); + return restTemplate.exchange( + baseURL + "/api/" + entityId.getEntityType() + "/" + entityId.getId() + "/calculatedFields?" + getUrlParams(pageLink), + HttpMethod.GET, HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, params).getBody(); + + } + private String getTimeUrlParams(TimePageLink pageLink) { String urlParams = getUrlParams(pageLink); if (pageLink.getStartTime() != null) {