@ -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 < String , List < Long > > gwLatencies = new HashMap < > ( ) ;
List < Long > transportLatencies = new ArrayList < > ( ) ;
publishLatency ( client , gwLatencies , transportLatencies , 5 ) ;
gatewayLatencyService . reportLatency ( ) ;
List < String > 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 < String , List < Map < String , Object > > > gatewayTelemetry = doGetAsyncTyped ( telemetryUrl , new TypeReference < > ( ) { } ) ;
Map < String , Object > latencyCheckTelemetry = gatewayTelemetry . get ( "latencyCheck" ) . get ( 0 ) ;
Map < String , GatewayLatencyState . ConnectorLatencyResult > 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 < String , List < Long > > gwLatencies , List < Long > 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 < GatewayLatencyData > 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 < Long > 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 < String > expectedKeys , byte [ ] payload , boolean withTs ) throws Exception {
processTelemetryTest ( topic , expectedKeys , payload , withTs , false ) ;
}