diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 785f124583..56dfe0ab46 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1020,6 +1020,8 @@ transport: # MQTT disconnect timeout in milliseconds. The time to wait for the client to disconnect after the server sends a disconnect message. disconnect_timeout: "${MQTT_DISCONNECT_TIMEOUT:1000}" msg_queue_size_per_device_limit: "${MQTT_MSG_QUEUE_SIZE_PER_DEVICE_LIMIT:100}" # messages await in the queue before the device connected state. This limit works on the low level before TenantProfileLimits mechanism + # Interval of periodic report of the gateway latency + gateway_latency_report_interval_sec: "${MQTT_GATEWAY_LATENCY_REPORT_INTERVAL_SEC:3600}" netty: # Netty leak detector level leak_detector_level: "${NETTY_LEAK_DETECTOR_LVL:DISABLED}" diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java index 7c8cdb5d6e..6deceb6fb9 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java @@ -18,30 +18,45 @@ package org.thingsboard.server.transport.mqtt.mqttv3.telemetry.timeseries; import com.fasterxml.jackson.core.type.TypeReference; import io.netty.handler.codec.mqtt.MqttQoS; import lombok.extern.slf4j.Slf4j; +import org.awaitility.Awaitility; import org.junit.Before; import org.junit.Test; +import org.mockito.ArgumentCaptor; +import org.springframework.boot.test.mock.mockito.SpyBean; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; +import org.thingsboard.server.transport.mqtt.gateway.GatewayLatencyService; +import org.thingsboard.server.transport.mqtt.gateway.latency.GatewayLatencyData; +import org.thingsboard.server.transport.mqtt.gateway.latency.GatewayLatencyState; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; +import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.OptionalLong; +import java.util.Random; import java.util.Set; import java.util.concurrent.TimeUnit; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.eq; +import static org.mockito.Mockito.verify; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_ATTRIBUTES_TOPIC; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_TELEMETRY_SHORT_JSON_TOPIC; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_TELEMETRY_SHORT_TOPIC; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_TELEMETRY_TOPIC; import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_CONNECT_TOPIC; +import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_LATENCY_TOPIC; import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_TELEMETRY_TOPIC; @Slf4j @@ -53,6 +68,9 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt protected static final String MALFORMED_JSON_PAYLOAD = "{\"key1\":, \"key2\":true, \"key3\": 3.0, \"key4\": 4," + " \"key5\": {\"someNumber\": 42, \"someArray\": [1,2,3], \"someNestedObject\": {\"key\": \"value\"}}}"; + @SpyBean + GatewayLatencyService gatewayLatencyService; + @Before public void beforeTest() throws Exception { MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() @@ -113,6 +131,86 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt client.disconnect(); } + @Test + public void testPushLatencyGateway() throws Exception { + MqttTestClient client = new MqttTestClient(); + client.connectAndWait(gatewayAccessToken); + + Map> gwLatencies = new HashMap<>(); + List transportLatencies = new ArrayList<>(); + + publishLatency(client, gwLatencies, transportLatencies, 5); + + gatewayLatencyService.reportLatency(); + + List actualKeys = getActualKeysList(savedGateway.getId(), List.of("latencyCheck")); + assertEquals("latencyCheck", actualKeys.get(0)); + + String telemetryUrl = String.format("/api/plugins/telemetry/DEVICE/%s/values/timeseries?startTs=%d&endTs=%d&keys=latencyCheck", savedGateway.getId(), 0, System.currentTimeMillis()); + + Map>> gatewayTelemetry = doGetAsyncTyped(telemetryUrl, new TypeReference<>() {}); + Map latencyCheckTelemetry = gatewayTelemetry.get("latencyCheck").get(0); + Map latencyCheckValue = JacksonUtil.fromString((String) latencyCheckTelemetry.get("value"), new TypeReference<>() {}); + assertNotNull(latencyCheckValue); + + long avgTransportLatency = (long) transportLatencies.stream().mapToLong(Long::longValue).average().getAsDouble(); + long minTransportLatency = transportLatencies.stream().mapToLong(Long::longValue).min().getAsLong(); + long maxTransportLatency = transportLatencies.stream().mapToLong(Long::longValue).max().getAsLong(); + + + gwLatencies.forEach((connectorName, gwLatencyList) -> { + long avgGwLatency = (long) gwLatencyList.stream().mapToLong(Long::longValue).average().getAsDouble(); + long minGwLatency = gwLatencyList.stream().mapToLong(Long::longValue).min().getAsLong(); + long maxGwLatency = gwLatencyList.stream().mapToLong(Long::longValue).max().getAsLong(); + + GatewayLatencyState.ConnectorLatencyResult connectorLatencyResult = latencyCheckValue.get(connectorName); + assertNotNull(connectorLatencyResult); + checkConnectorLatencyResult(connectorLatencyResult, avgGwLatency, minGwLatency, maxGwLatency, avgTransportLatency, minTransportLatency, maxTransportLatency); + }); + + client.disconnect(); + + Awaitility.await() + .atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> verify(gatewayLatencyService).onDeviceDisconnect(savedGateway.getId())); + } + + private void publishLatency(MqttTestClient client, Map> gwLatencies, List transportLatencies, int n) throws Exception { + Random random = new Random(); + for (int i = 0; i < n; i++) { + long publishedTs = System.currentTimeMillis(); + long gatewayLatencyA = random.nextLong(1000, 5000); + long gatewayLatencyB = random.nextLong(1200, 4500); + long transportReceiveTs = publishLatencyAndGetTransportReceiveTs(client, publishedTs, gatewayLatencyA, gatewayLatencyB); + gwLatencies.computeIfAbsent("connectorA", key -> new ArrayList<>()).add(gatewayLatencyA); + gwLatencies.computeIfAbsent("connectorB", key -> new ArrayList<>()).add(gatewayLatencyB); + transportLatencies.add(transportReceiveTs - publishedTs); + Thread.sleep(1); + } + } + + private long publishLatencyAndGetTransportReceiveTs(MqttTestClient client, long publishedTs, long gatewayLatencyA, long gatewayLatencyB) throws Exception { + List data = new ArrayList<>(); + data.add(new GatewayLatencyData("connectorA", publishedTs - gatewayLatencyA, publishedTs)); + data.add(new GatewayLatencyData("connectorB", publishedTs - gatewayLatencyB, publishedTs)); + + client.publishAndWait(GATEWAY_LATENCY_TOPIC, JacksonUtil.writeValueAsBytes(data)); + ArgumentCaptor transportReceiveTsCaptor = ArgumentCaptor.forClass(Long.class); + verify(gatewayLatencyService).process(any(), eq(savedGateway.getId()), eq(data), transportReceiveTsCaptor.capture()); + return transportReceiveTsCaptor.getValue(); + } + + private void checkConnectorLatencyResult(GatewayLatencyState.ConnectorLatencyResult result, long avgGwLatency, long minGwLatency, long maxGwLatency, + long avgTransportLatency, long minTransportLatency, long maxTransportLatency) { + assertNotNull(result); + assertEquals(avgGwLatency, result.avgGwLatency()); + assertEquals(minGwLatency, result.minGwLatency()); + assertEquals(maxGwLatency, result.maxGwLatency()); + assertEquals(avgTransportLatency, result.transportLatencyAvg()); + assertEquals(minTransportLatency, result.minTransportLatency()); + assertEquals(maxTransportLatency, result.maxTransportLatency()); + } + protected void processJsonPayloadTelemetryTest(String topic, List expectedKeys, byte[] payload, boolean withTs) throws Exception { processTelemetryTest(topic, expectedKeys, payload, withTs, false); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java index cabdea145c..8945f2d084 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttTopics.java @@ -75,6 +75,7 @@ public class MqttTopics { public static final String GATEWAY_RPC_TOPIC = BASE_GATEWAY_API_TOPIC + RPC; public static final String GATEWAY_ATTRIBUTES_REQUEST_TOPIC = BASE_GATEWAY_API_TOPIC + ATTRIBUTES_REQUEST; public static final String GATEWAY_ATTRIBUTES_RESPONSE_TOPIC = BASE_GATEWAY_API_TOPIC + ATTRIBUTES_RESPONSE; + public static final String GATEWAY_LATENCY_TOPIC = BASE_GATEWAY_API_TOPIC + "/latency"; // v2 topics public static final String BASE_DEVICE_API_TOPIC_V2 = "v2"; public static final String REQUEST_ID_PATTERN = "(?\\d+)"; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 16d72ff89f..555d7de868 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -22,11 +22,11 @@ import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; import org.thingsboard.server.common.transport.TransportContext; import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor; +import org.thingsboard.server.transport.mqtt.gateway.GatewayLatencyService; import java.net.InetSocketAddress; import java.util.concurrent.atomic.AtomicInteger; @@ -36,7 +36,7 @@ import java.util.concurrent.atomic.AtomicInteger; */ @Slf4j @Component -@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')") +@TbMqttTransportComponent public class MqttTransportContext extends TransportContext { @Getter @@ -51,6 +51,10 @@ public class MqttTransportContext extends TransportContext { @Autowired private ProtoMqttAdaptor protoMqttAdaptor; + @Getter + @Autowired + private GatewayLatencyService gatewayLatencyService; + @Getter @Value("${transport.mqtt.netty.max_payload_size}") private Integer maxPayloadSize; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index eeeb287dc8..ef63fb5902 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -419,6 +419,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement case MqttTopics.GATEWAY_DISCONNECT_TOPIC: gatewaySessionHandler.onDeviceDisconnect(mqttMsg); break; + case MqttTopics.GATEWAY_LATENCY_TOPIC: + gatewaySessionHandler.onGatewayLatency(mqttMsg); + break; default: ack(ctx, msgId, MqttReasonCodes.PubAck.TOPIC_NAME_INVALID); } @@ -1187,7 +1190,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(deviceSessionCtx.getSessionInfo(), SESSION_EVENT_MSG_CLOSED, null); transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); if (gatewaySessionHandler != null) { - gatewaySessionHandler.onDevicesDisconnect(); + gatewaySessionHandler.onGatewayDisconnect(); } if (sparkplugSessionHandler != null) { // add Msg Telemetry node: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() @@ -1418,7 +1421,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional deviceProfileOpt) { deviceSessionCtx.onDeviceUpdate(sessionInfo, device, deviceProfileOpt); if (gatewaySessionHandler != null) { - gatewaySessionHandler.onDeviceUpdate(sessionInfo, device, deviceProfileOpt); + gatewaySessionHandler.onGatewayUpdate(sessionInfo, device, deviceProfileOpt); } } @@ -1427,6 +1430,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement context.onAuthFailure(address); ChannelHandlerContext ctx = deviceSessionCtx.getChannel(); closeCtx(ctx, MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION); + if (gatewaySessionHandler != null) { + gatewaySessionHandler.onGatewayDelete(deviceId); + } } public void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, ThingsboardErrorCode result, String errorMsg) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index ffc7e9e561..b8d1dbe582 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java @@ -39,7 +39,7 @@ import java.net.InetSocketAddress; * @author Andrew Shvayka */ @Service("MqttTransportService") -@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')") +@TbMqttTransportComponent @Slf4j public class MqttTransportService implements TbTransportService { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TbMqttTransportComponent.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TbMqttTransportComponent.java new file mode 100644 index 0000000000..da0047cc7e --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/TbMqttTransportComponent.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2024 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.server.transport.mqtt; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; + +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; + +@Retention(RetentionPolicy.RUNTIME) +@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')") +public @interface TbMqttTransportComponent { +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/GatewayLatencyService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/GatewayLatencyService.java new file mode 100644 index 0000000000..c6e9002877 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/GatewayLatencyService.java @@ -0,0 +1,123 @@ +/** + * Copyright © 2016-2024 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.server.transport.mqtt.gateway; + +import jakarta.annotation.PostConstruct; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.common.transport.TransportServiceCallback; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.scheduler.SchedulerComponent; +import org.thingsboard.server.transport.mqtt.TbMqttTransportComponent; +import org.thingsboard.server.transport.mqtt.gateway.latency.GatewayLatencyData; +import org.thingsboard.server.transport.mqtt.gateway.latency.GatewayLatencyState; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; + +@Slf4j +@Service +@TbMqttTransportComponent +public class GatewayLatencyService { + + @Value("${transport.mqtt.gateway_latency_report_interval_sec:3600}") + private int latencyReportIntervalSec; + + @Autowired + private SchedulerComponent scheduler; + + @Autowired + private TransportService transportService; + + private Map states = new ConcurrentHashMap<>(); + + @PostConstruct + private void init() { + scheduler.scheduleAtFixedRate(this::reportLatency, latencyReportIntervalSec, latencyReportIntervalSec, TimeUnit.SECONDS); + } + + public void process(TransportProtos.SessionInfoProto sessionInfo, DeviceId gatewayId, List data, long ts) { + states.computeIfAbsent(gatewayId, k -> new GatewayLatencyState(sessionInfo)).update(ts, data); + } + + public void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceId gatewayId) { + var state = states.get(gatewayId); + if (state != null) { + state.updateSessionInfo(sessionInfo); + } + } + + public void onDeviceDelete(DeviceId deviceId) { + var state = states.remove(deviceId); + if (state != null) { + state.clear(); + } + } + + public void onDeviceDisconnect(DeviceId deviceId) { + GatewayLatencyState state = states.remove(deviceId); + if (state != null) { + reportLatency(state, System.currentTimeMillis()); + } + } + + public void reportLatency() { + if (states.isEmpty()) { + return; + } + Map oldStates = states; + states = new ConcurrentHashMap<>(); + + long ts = System.currentTimeMillis(); + + oldStates.forEach((gatewayId, state) -> { + reportLatency(state, ts); + }); + oldStates.clear(); + } + + private void reportLatency(GatewayLatencyState state, long ts) { + if (state.isEmpty()) { + return; + } + var result = state.getLatencyStateResult(); + state.clear(); + var kvProto = TransportProtos.KeyValueProto.newBuilder() + .setKey("latencyCheck") + .setType(TransportProtos.KeyValueType.JSON_V) + .setJsonV(JacksonUtil.toString(result)) + .build(); + + TransportProtos.TsKvListProto tsKvList = TransportProtos.TsKvListProto.newBuilder() + .setTs(ts) + .addKv(kvProto) + .build(); + + TransportProtos.PostTelemetryMsg telemetryMsg = TransportProtos.PostTelemetryMsg.newBuilder() + .addTsKvList(tsKvList) + .build(); + + transportService.process(state.getSessionInfo(), telemetryMsg, TransportServiceCallback.EMPTY); + } + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyData.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyData.java new file mode 100644 index 0000000000..6260c8356e --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyData.java @@ -0,0 +1,19 @@ +/** + * Copyright © 2016-2024 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.server.transport.mqtt.gateway.latency; + +public record GatewayLatencyData(String connectorName, long receivedTs, long publishedTs) { +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyState.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyState.java new file mode 100644 index 0000000000..3ed902f81d --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/gateway/latency/GatewayLatencyState.java @@ -0,0 +1,119 @@ +/** + * Copyright © 2016-2024 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.server.transport.mqtt.gateway.latency; + +import lombok.Getter; +import org.thingsboard.server.gen.transport.TransportProtos; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + +public class GatewayLatencyState { + + private final Map connectors; + private final Lock updateLock; + + @Getter + private volatile TransportProtos.SessionInfoProto sessionInfo; + + public GatewayLatencyState(TransportProtos.SessionInfoProto sessionInfo) { + this.connectors = new HashMap<>(); + this.updateLock = new ReentrantLock(); + this.sessionInfo = sessionInfo; + } + + public void updateSessionInfo(TransportProtos.SessionInfoProto sessionInfo) { + this.sessionInfo = sessionInfo; + } + + public void update(long ts, List latencyData) { + updateLock.lock(); + try { + latencyData.forEach(data -> { + connectors.computeIfAbsent(data.connectorName(), k -> new ConnectorLatencyState()).update(ts, data); + }); + } finally { + updateLock.unlock(); + } + } + + public void clear() { + connectors.clear(); + } + + public Map getLatencyStateResult() { + Map result = new HashMap<>(); + connectors.forEach((name, state) -> result.put(name, state.getResult())); + return result; + } + + public boolean isEmpty() { + return connectors.isEmpty(); + } + + private static class ConnectorLatencyState { + private final AtomicInteger count; + private final AtomicLong gwLatencySum; + private final AtomicLong transportLatencySum; + private volatile long minGwLatency; + private volatile long maxGwLatency; + private volatile long minTransportLatency; + private volatile long maxTransportLatency; + + private ConnectorLatencyState() { + this.count = new AtomicInteger(0); + this.gwLatencySum = new AtomicLong(0); + this.transportLatencySum = new AtomicLong(0); + } + + private void update(long serverReceiveTs, GatewayLatencyData latencyData) { + long gwLatency = latencyData.publishedTs() - latencyData.receivedTs(); + long transportLatency = serverReceiveTs - latencyData.publishedTs(); + count.incrementAndGet(); + gwLatencySum.addAndGet(gwLatency); + transportLatencySum.addAndGet(transportLatency); + if (minGwLatency == 0 || minGwLatency > gwLatency) { + minGwLatency = gwLatency; + } + if (maxGwLatency < gwLatency) { + maxGwLatency = gwLatency; + } + if (minTransportLatency == 0 || minTransportLatency > transportLatency) { + minTransportLatency = transportLatency; + } + if (maxTransportLatency < transportLatency) { + maxTransportLatency = transportLatency; + } + } + + private ConnectorLatencyResult getResult() { + long count = this.count.get(); + long avgGwLatency = gwLatencySum.get() / count; + long transportLatencyAvg = transportLatencySum.get() / count; + return new ConnectorLatencyResult(avgGwLatency, minGwLatency, maxGwLatency, transportLatencyAvg, minTransportLatency, maxTransportLatency); + } + } + + public record ConnectorLatencyResult(long avgGwLatency, long minGwLatency, long maxGwLatency, + long transportLatencyAvg, long minTransportLatency, long maxTransportLatency) { + } + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java index 9aa228afff..684f6b3114 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java @@ -15,20 +15,36 @@ */ package org.thingsboard.server.transport.mqtt.session; +import com.fasterxml.jackson.core.type.TypeReference; import io.netty.buffer.ByteBuf; import io.netty.handler.codec.mqtt.MqttPublishMessage; +import io.netty.handler.codec.mqtt.MqttReasonCodes; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.adaptor.AdaptorException; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.mqtt.gateway.GatewayLatencyService; +import java.util.Optional; import java.util.UUID; +import static com.amazonaws.util.StringUtils.UTF8; + /** * Created by nickAS21 on 26.12.22 */ +@Slf4j public class GatewaySessionHandler extends AbstractGatewaySessionHandler { + private final GatewayLatencyService latencyService; + public GatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId, boolean overwriteDevicesActivity) { super(deviceSessionCtx, sessionId, overwriteDevicesActivity); + this.latencyService = deviceSessionCtx.getContext().getGatewayLatencyService(); } public void onDeviceConnect(MqttPublishMessage mqttMsg) throws AdaptorException { @@ -51,7 +67,38 @@ public class GatewaySessionHandler extends AbstractGatewaySessionHandler deviceProfileOpt) { + this.onDeviceUpdate(sessionInfo, device, deviceProfileOpt); + latencyService.onDeviceUpdate(sessionInfo, gateway.getDeviceId()); + } + + public void onGatewayDelete(DeviceId deviceId) { + latencyService.onDeviceDelete(deviceId); + } + + public void onGatewayLatency(MqttPublishMessage mqttMsg) throws AdaptorException { + int msgId = getMsgId(mqttMsg); + ByteBuf payloadData = mqttMsg.payload(); + String payload = payloadData.toString(UTF8); + if (payload == null) { + log.debug("[{}][{}][{}] Payload is empty!", gateway.getTenantId(), gateway.getDeviceId(), sessionId); + throw new AdaptorException(new IllegalArgumentException("Payload is empty!")); + } + long ts = System.currentTimeMillis(); + try { + latencyService.process(deviceSessionCtx.getSessionInfo(), gateway.getDeviceId(), JacksonUtil.fromString(payload, new TypeReference<>() {}), ts); + ack(msgId, MqttReasonCodes.PubAck.SUCCESS); + } catch (IllegalArgumentException e) { + throw new AdaptorException(e); + } } }