Browse Source
# Conflicts: # application/src/main/data/upgrade/3.1.0/schema_update.sqlpull/3557/head
244 changed files with 16905 additions and 18002 deletions
File diff suppressed because one or more lines are too long
@ -0,0 +1,35 @@ |
|||
-- |
|||
-- Copyright © 2016-2020 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. |
|||
-- |
|||
|
|||
CREATE TABLE IF NOT EXISTS ts_kv_latest |
|||
( |
|||
entity_id uuid NOT NULL, |
|||
key int NOT NULL, |
|||
ts bigint NOT NULL, |
|||
bool_v boolean, |
|||
str_v varchar(10000000), |
|||
long_v bigint, |
|||
dbl_v double precision, |
|||
json_v json, |
|||
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_id, key) |
|||
); |
|||
|
|||
CREATE TABLE IF NOT EXISTS ts_kv_dictionary |
|||
( |
|||
key varchar(255) NOT NULL, |
|||
key_id serial UNIQUE, |
|||
CONSTRAINT ts_key_id_pkey PRIMARY KEY (key) |
|||
); |
|||
@ -0,0 +1,40 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.query; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Component; |
|||
|
|||
import java.util.Arrays; |
|||
|
|||
@Component |
|||
@Slf4j |
|||
public class DefaultQueryLogComponent implements QueryLogComponent { |
|||
|
|||
@Value("${sql.log_queries:false}") |
|||
private boolean logSqlQueries; |
|||
@Value("${sql.log_queries_threshold:5000}") |
|||
private long logQueriesThreshold; |
|||
|
|||
@Override |
|||
public void logQuery(QueryContext ctx, String query, long duration) { |
|||
if (logSqlQueries && duration > logQueriesThreshold) { |
|||
log.info("QUERY: {} took {}ms", query, duration); |
|||
Arrays.asList(ctx.getParameterNames()).forEach(param -> log.info("QUERY PARAM: {} -> {}", param, ctx.getValue(param))); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.query; |
|||
|
|||
public interface QueryLogComponent { |
|||
|
|||
void logQuery(QueryContext ctx, String query, long duration); |
|||
} |
|||
@ -0,0 +1,40 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.util; |
|||
|
|||
import org.springframework.jdbc.core.JdbcTemplate; |
|||
|
|||
import java.sql.DriverManager; |
|||
|
|||
public class DaoTestUtil { |
|||
private static final String POSTGRES_DRIVER_CLASS = "org.postgresql.Driver"; |
|||
private static final String H2_DRIVER_CLASS = "org.hsqldb.jdbc.JDBCDriver"; |
|||
|
|||
|
|||
public static SqlDbType getSqlDbType(JdbcTemplate template){ |
|||
try { |
|||
String driverName = DriverManager.getDriver(template.getDataSource().getConnection().getMetaData().getURL()).getClass().getName(); |
|||
if (POSTGRES_DRIVER_CLASS.equals(driverName)) { |
|||
return SqlDbType.POSTGRES; |
|||
} else if (H2_DRIVER_CLASS.equals(driverName)) { |
|||
return SqlDbType.H2; |
|||
} |
|||
} catch (Exception e) { |
|||
e.printStackTrace(); |
|||
} |
|||
return null; |
|||
} |
|||
} |
|||
@ -0,0 +1,20 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.util; |
|||
|
|||
public enum SqlDbType { |
|||
POSTGRES, H2; |
|||
} |
|||
@ -0,0 +1,57 @@ |
|||
# |
|||
# Copyright © 2016-2020 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. |
|||
# |
|||
|
|||
version: '2.2' |
|||
|
|||
services: |
|||
tb-js-executor: |
|||
env_file: |
|||
- queue-confluent.env |
|||
tb-core1: |
|||
env_file: |
|||
- queue-confluent.env |
|||
depends_on: |
|||
- redis |
|||
tb-core2: |
|||
env_file: |
|||
- queue-confluent.env |
|||
depends_on: |
|||
- redis |
|||
tb-rule-engine1: |
|||
env_file: |
|||
- queue-confluent.env |
|||
depends_on: |
|||
- redis |
|||
tb-rule-engine2: |
|||
env_file: |
|||
- queue-confluent.env |
|||
depends_on: |
|||
- redis |
|||
tb-mqtt-transport1: |
|||
env_file: |
|||
- queue-confluent.env |
|||
tb-mqtt-transport2: |
|||
env_file: |
|||
- queue-confluent.env |
|||
tb-http-transport1: |
|||
env_file: |
|||
- queue-confluent.env |
|||
tb-http-transport2: |
|||
env_file: |
|||
- queue-confluent.env |
|||
tb-coap-transport: |
|||
env_file: |
|||
- queue-confluent.env |
|||
@ -0,0 +1,18 @@ |
|||
TB_QUEUE_TYPE=kafka |
|||
|
|||
TB_KAFKA_SERVERS=confluent.cloud:9092 |
|||
TB_QUEUE_KAFKA_REPLICATION_FACTOR=3 |
|||
|
|||
TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD=true |
|||
TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM=https |
|||
TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM=PLAIN |
|||
TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG=org.apache.kafka.common.security.plain.PlainLoginModule required username="CLUSTER_API_KEY" password="CLUSTER_API_SECRET"; |
|||
TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL=SASL_SSL |
|||
TB_QUEUE_KAFKA_CONFLUENT_USERNAME=CLUSTER_API_KEY |
|||
TB_QUEUE_KAFKA_CONFLUENT_PASSWORD=CLUSTER_API_SECRET |
|||
|
|||
TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 |
|||
TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 |
|||
TB_QUEUE_KAFKA_TA_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 |
|||
TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 |
|||
TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:104857600 |
|||
@ -0,0 +1,376 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.msa.connectivity; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import com.google.common.collect.Sets; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.ListeningExecutorService; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import com.google.gson.JsonObject; |
|||
import com.google.gson.JsonParser; |
|||
import io.netty.buffer.ByteBuf; |
|||
import io.netty.buffer.Unpooled; |
|||
import io.netty.handler.codec.mqtt.MqttQoS; |
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.junit.After; |
|||
import org.junit.Assert; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.springframework.http.ResponseEntity; |
|||
import org.thingsboard.mqtt.MqttClient; |
|||
import org.thingsboard.mqtt.MqttClientConfig; |
|||
import org.thingsboard.mqtt.MqttHandler; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|||
import org.thingsboard.server.common.data.security.DeviceCredentials; |
|||
import org.thingsboard.server.msa.AbstractContainerTest; |
|||
import org.thingsboard.server.msa.WsClient; |
|||
import org.thingsboard.server.msa.mapper.WsTelemetryResponse; |
|||
|
|||
import java.io.IOException; |
|||
import java.nio.charset.StandardCharsets; |
|||
import java.util.List; |
|||
import java.util.Objects; |
|||
import java.util.Optional; |
|||
import java.util.Random; |
|||
import java.util.concurrent.ArrayBlockingQueue; |
|||
import java.util.concurrent.BlockingQueue; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@Slf4j |
|||
public class MqttGatewayClientTest extends AbstractContainerTest { |
|||
Device gatewayDevice; |
|||
MqttClient mqttClient; |
|||
Device createdDevice; |
|||
MqttMessageListener listener; |
|||
|
|||
@Before |
|||
public void createGateway() throws Exception { |
|||
restClient.login("tenant@thingsboard.org", "tenant"); |
|||
this.gatewayDevice = createGatewayDevice(); |
|||
Optional<DeviceCredentials> gatewayDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(gatewayDevice.getId()); |
|||
Assert.assertTrue(gatewayDeviceCredentials.isPresent()); |
|||
this.listener = new MqttMessageListener(); |
|||
this.mqttClient = getMqttClient(gatewayDeviceCredentials.get(), listener); |
|||
this.createdDevice = createDeviceThroughGateway(mqttClient, gatewayDevice); |
|||
} |
|||
|
|||
@After |
|||
public void removeGateway() throws Exception { |
|||
restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + this.gatewayDevice.getId()); |
|||
restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + this.createdDevice.getId()); |
|||
this.listener = null; |
|||
this.mqttClient = null; |
|||
this.createdDevice = null; |
|||
} |
|||
|
|||
@Test |
|||
public void telemetryUpload() throws Exception { |
|||
WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); |
|||
mqttClient.publish("v1/gateway/telemetry", Unpooled.wrappedBuffer(createGatewayPayload(createdDevice.getName(), -1).toString().getBytes())).get(); |
|||
WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); |
|||
log.info("Received telemetry: {}", actualLatestTelemetry); |
|||
wsClient.closeBlocking(); |
|||
|
|||
Assert.assertEquals(4, actualLatestTelemetry.getData().size()); |
|||
Assert.assertEquals(Sets.newHashSet("booleanKey", "stringKey", "doubleKey", "longKey"), |
|||
actualLatestTelemetry.getLatestValues().keySet()); |
|||
|
|||
Assert.assertTrue(verify(actualLatestTelemetry, "booleanKey", Boolean.TRUE.toString())); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "stringKey", "value1")); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "doubleKey", Double.toString(42.0))); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "longKey", Long.toString(73))); |
|||
} |
|||
|
|||
@Test |
|||
public void telemetryUploadWithTs() throws Exception { |
|||
long ts = 1451649600512L; |
|||
|
|||
restClient.login("tenant@thingsboard.org", "tenant"); |
|||
WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); |
|||
mqttClient.publish("v1/gateway/telemetry", Unpooled.wrappedBuffer(createGatewayPayload(createdDevice.getName(), ts).toString().getBytes())).get(); |
|||
WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); |
|||
log.info("Received telemetry: {}", actualLatestTelemetry); |
|||
wsClient.closeBlocking(); |
|||
|
|||
Assert.assertEquals(4, actualLatestTelemetry.getData().size()); |
|||
Assert.assertEquals(getExpectedLatestValues(ts), actualLatestTelemetry.getLatestValues()); |
|||
|
|||
Assert.assertTrue(verify(actualLatestTelemetry, "booleanKey", ts, Boolean.TRUE.toString())); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "stringKey", ts, "value1")); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "doubleKey", ts, Double.toString(42.0))); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "longKey", ts, Long.toString(73))); |
|||
} |
|||
|
|||
@Test |
|||
public void publishAttributeUpdateToServer() throws Exception { |
|||
Optional<DeviceCredentials> createdDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); |
|||
Assert.assertTrue(createdDeviceCredentials.isPresent()); |
|||
WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); |
|||
JsonObject clientAttributes = new JsonObject(); |
|||
clientAttributes.addProperty("attr1", "value1"); |
|||
clientAttributes.addProperty("attr2", true); |
|||
clientAttributes.addProperty("attr3", 42.0); |
|||
clientAttributes.addProperty("attr4", 73); |
|||
JsonObject gatewayClientAttributes = new JsonObject(); |
|||
gatewayClientAttributes.add(createdDevice.getName(), clientAttributes); |
|||
mqttClient.publish("v1/gateway/attributes", Unpooled.wrappedBuffer(gatewayClientAttributes.toString().getBytes())).get(); |
|||
WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); |
|||
log.info("Received attributes: {}", actualLatestTelemetry); |
|||
wsClient.closeBlocking(); |
|||
|
|||
Assert.assertEquals(4, actualLatestTelemetry.getData().size()); |
|||
Assert.assertEquals(Sets.newHashSet("attr1", "attr2", "attr3", "attr4"), |
|||
actualLatestTelemetry.getLatestValues().keySet()); |
|||
|
|||
Assert.assertTrue(verify(actualLatestTelemetry, "attr1", "value1")); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "attr2", Boolean.TRUE.toString())); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "attr3", Double.toString(42.0))); |
|||
Assert.assertTrue(verify(actualLatestTelemetry, "attr4", Long.toString(73))); |
|||
} |
|||
|
|||
@Test |
|||
public void requestAttributeValuesFromServer() throws Exception { |
|||
WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); |
|||
// Add a new client attribute
|
|||
JsonObject clientAttributes = new JsonObject(); |
|||
String clientAttributeValue = RandomStringUtils.randomAlphanumeric(8); |
|||
clientAttributes.addProperty("clientAttr", clientAttributeValue); |
|||
|
|||
JsonObject gatewayClientAttributes = new JsonObject(); |
|||
gatewayClientAttributes.add(createdDevice.getName(), clientAttributes); |
|||
mqttClient.publish("v1/gateway/attributes", Unpooled.wrappedBuffer(gatewayClientAttributes.toString().getBytes())).get(); |
|||
|
|||
WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); |
|||
log.info("Received ws telemetry: {}", actualLatestTelemetry); |
|||
wsClient.closeBlocking(); |
|||
|
|||
Assert.assertEquals(1, actualLatestTelemetry.getData().size()); |
|||
Assert.assertEquals(Sets.newHashSet("clientAttr"), |
|||
actualLatestTelemetry.getLatestValues().keySet()); |
|||
|
|||
Assert.assertTrue(verify(actualLatestTelemetry, "clientAttr", clientAttributeValue)); |
|||
|
|||
// Add a new shared attribute
|
|||
JsonObject sharedAttributes = new JsonObject(); |
|||
String sharedAttributeValue = RandomStringUtils.randomAlphanumeric(8); |
|||
sharedAttributes.addProperty("sharedAttr", sharedAttributeValue); |
|||
|
|||
ResponseEntity sharedAttributesResponse = restClient.getRestTemplate() |
|||
.postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", |
|||
mapper.readTree(sharedAttributes.toString()), ResponseEntity.class, |
|||
createdDevice.getId()); |
|||
Assert.assertTrue(sharedAttributesResponse.getStatusCode().is2xxSuccessful()); |
|||
MqttEvent sharedAttributeEvent = listener.getEvents().poll(10, TimeUnit.SECONDS); |
|||
|
|||
// Catch attribute update event
|
|||
Assert.assertNotNull(sharedAttributeEvent); |
|||
Assert.assertEquals("v1/gateway/attributes", sharedAttributeEvent.getTopic()); |
|||
|
|||
// Subscribe to attributes response
|
|||
mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); |
|||
|
|||
// Wait until subscription is processed
|
|||
TimeUnit.SECONDS.sleep(3); |
|||
|
|||
checkAttribute(true, clientAttributeValue); |
|||
checkAttribute(false, sharedAttributeValue); |
|||
} |
|||
|
|||
@Test |
|||
public void subscribeToAttributeUpdatesFromServer() throws Exception { |
|||
mqttClient.on("v1/gateway/attributes", listener, MqttQoS.AT_LEAST_ONCE).get(); |
|||
// Wait until subscription is processed
|
|||
TimeUnit.SECONDS.sleep(3); |
|||
String sharedAttributeName = "sharedAttr"; |
|||
// Add a new shared attribute
|
|||
|
|||
JsonObject sharedAttributes = new JsonObject(); |
|||
String sharedAttributeValue = RandomStringUtils.randomAlphanumeric(8); |
|||
sharedAttributes.addProperty(sharedAttributeName, sharedAttributeValue); |
|||
|
|||
JsonObject gatewaySharedAttributeValue = new JsonObject(); |
|||
gatewaySharedAttributeValue.addProperty("device", createdDevice.getName()); |
|||
gatewaySharedAttributeValue.add("data", sharedAttributes); |
|||
|
|||
ResponseEntity sharedAttributesResponse = restClient.getRestTemplate() |
|||
.postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", |
|||
mapper.readTree(sharedAttributes.toString()), ResponseEntity.class, |
|||
createdDevice.getId()); |
|||
Assert.assertTrue(sharedAttributesResponse.getStatusCode().is2xxSuccessful()); |
|||
|
|||
MqttEvent event = listener.getEvents().poll(10, TimeUnit.SECONDS); |
|||
Assert.assertEquals(sharedAttributeValue, |
|||
mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()); |
|||
|
|||
// Update the shared attribute value
|
|||
JsonObject updatedSharedAttributes = new JsonObject(); |
|||
String updatedSharedAttributeValue = RandomStringUtils.randomAlphanumeric(8); |
|||
updatedSharedAttributes.addProperty(sharedAttributeName, updatedSharedAttributeValue); |
|||
|
|||
JsonObject gatewayUpdatedSharedAttributeValue = new JsonObject(); |
|||
gatewayUpdatedSharedAttributeValue.addProperty("device", createdDevice.getName()); |
|||
gatewayUpdatedSharedAttributeValue.add("data", updatedSharedAttributes); |
|||
|
|||
ResponseEntity updatedSharedAttributesResponse = restClient.getRestTemplate() |
|||
.postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", |
|||
mapper.readTree(updatedSharedAttributes.toString()), ResponseEntity.class, |
|||
createdDevice.getId()); |
|||
Assert.assertTrue(updatedSharedAttributesResponse.getStatusCode().is2xxSuccessful()); |
|||
|
|||
event = listener.getEvents().poll(10, TimeUnit.SECONDS); |
|||
Assert.assertEquals(updatedSharedAttributeValue, |
|||
mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()); |
|||
} |
|||
|
|||
@Test |
|||
public void serverSideRpc() throws Exception { |
|||
String gatewayRpcTopic = "v1/gateway/rpc"; |
|||
mqttClient.on(gatewayRpcTopic, listener, MqttQoS.AT_LEAST_ONCE).get(); |
|||
|
|||
// Wait until subscription is processed
|
|||
TimeUnit.SECONDS.sleep(3); |
|||
|
|||
// Send an RPC from the server
|
|||
JsonObject serverRpcPayload = new JsonObject(); |
|||
serverRpcPayload.addProperty("method", "getValue"); |
|||
serverRpcPayload.addProperty("params", true); |
|||
ListeningExecutorService service = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); |
|||
ListenableFuture<ResponseEntity> future = service.submit(() -> { |
|||
try { |
|||
return restClient.getRestTemplate() |
|||
.postForEntity(HTTPS_URL + "/api/plugins/rpc/twoway/{deviceId}", |
|||
mapper.readTree(serverRpcPayload.toString()), String.class, |
|||
createdDevice.getId()); |
|||
} catch (IOException e) { |
|||
return ResponseEntity.badRequest().build(); |
|||
} |
|||
}); |
|||
|
|||
// Wait for RPC call from the server and send the response
|
|||
MqttEvent requestFromServer = listener.getEvents().poll(10, TimeUnit.SECONDS); |
|||
|
|||
Assert.assertNotNull(requestFromServer); |
|||
Assert.assertNotNull(requestFromServer.getMessage()); |
|||
|
|||
JsonObject requestFromServerJson = new JsonParser().parse(requestFromServer.getMessage()).getAsJsonObject(); |
|||
|
|||
Assert.assertEquals(createdDevice.getName(), requestFromServerJson.get("device").getAsString()); |
|||
|
|||
JsonObject requestFromServerData = requestFromServerJson.get("data").getAsJsonObject(); |
|||
|
|||
Assert.assertEquals("getValue", requestFromServerData.get("method").getAsString()); |
|||
Assert.assertTrue(requestFromServerData.get("params").getAsBoolean()); |
|||
|
|||
int requestId = requestFromServerData.get("id").getAsInt(); |
|||
|
|||
JsonObject clientResponse = new JsonObject(); |
|||
clientResponse.addProperty("response", "someResponse"); |
|||
JsonObject gatewayResponse = new JsonObject(); |
|||
gatewayResponse.addProperty("device", createdDevice.getName()); |
|||
gatewayResponse.addProperty("id", requestId); |
|||
gatewayResponse.add("data", clientResponse); |
|||
// Send a response to the server's RPC request
|
|||
|
|||
mqttClient.publish(gatewayRpcTopic, Unpooled.wrappedBuffer(gatewayResponse.toString().getBytes())).get(); |
|||
ResponseEntity serverResponse = future.get(5, TimeUnit.SECONDS); |
|||
Assert.assertTrue(serverResponse.getStatusCode().is2xxSuccessful()); |
|||
Assert.assertEquals(clientResponse.toString(), serverResponse.getBody()); |
|||
} |
|||
|
|||
private void checkAttribute(boolean client, String expectedValue) throws Exception{ |
|||
JsonObject gatewayAttributesRequest = new JsonObject(); |
|||
int messageId = new Random().nextInt(100); |
|||
gatewayAttributesRequest.addProperty("id", messageId); |
|||
gatewayAttributesRequest.addProperty("device", createdDevice.getName()); |
|||
gatewayAttributesRequest.addProperty("client", client); |
|||
String attributeName; |
|||
if (client) |
|||
attributeName = "clientAttr"; |
|||
else |
|||
attributeName = "sharedAttr"; |
|||
gatewayAttributesRequest.addProperty("key", attributeName); |
|||
log.info(gatewayAttributesRequest.toString()); |
|||
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(gatewayAttributesRequest.toString().getBytes())).get(); |
|||
MqttEvent clientAttributeEvent = listener.getEvents().poll(10, TimeUnit.SECONDS); |
|||
Assert.assertNotNull(clientAttributeEvent); |
|||
JsonObject responseMessage = new JsonParser().parse(Objects.requireNonNull(clientAttributeEvent).getMessage()).getAsJsonObject(); |
|||
|
|||
Assert.assertEquals(messageId, responseMessage.get("id").getAsInt()); |
|||
Assert.assertEquals(createdDevice.getName(), responseMessage.get("device").getAsString()); |
|||
Assert.assertEquals(3, responseMessage.entrySet().size()); |
|||
Assert.assertEquals(expectedValue, responseMessage.get("value").getAsString()); |
|||
} |
|||
|
|||
private Device createDeviceThroughGateway(MqttClient mqttClient, Device gatewayDevice) throws Exception { |
|||
String deviceName = "mqtt_device"; |
|||
mqttClient.publish("v1/gateway/connect", Unpooled.wrappedBuffer(createGatewayConnectPayload(deviceName).toString().getBytes())).get(); |
|||
|
|||
TimeUnit.SECONDS.sleep(3); |
|||
List<EntityRelation> relations = restClient.findByFrom(gatewayDevice.getId(), RelationTypeGroup.COMMON); |
|||
|
|||
Assert.assertEquals(1, relations.size()); |
|||
|
|||
EntityId createdEntityId = relations.get(0).getTo(); |
|||
DeviceId createdDeviceId = new DeviceId(createdEntityId.getId()); |
|||
Optional<Device> createdDevice = restClient.getDeviceById(createdDeviceId); |
|||
|
|||
Assert.assertTrue(createdDevice.isPresent()); |
|||
|
|||
return createdDevice.get(); |
|||
} |
|||
|
|||
private MqttClient getMqttClient(DeviceCredentials deviceCredentials, MqttMessageListener listener) throws InterruptedException, ExecutionException { |
|||
MqttClientConfig clientConfig = new MqttClientConfig(); |
|||
clientConfig.setClientId("MQTT client from test"); |
|||
clientConfig.setUsername(deviceCredentials.getCredentialsId()); |
|||
MqttClient mqttClient = MqttClient.create(clientConfig, listener); |
|||
mqttClient.connect("localhost", 1883).get(); |
|||
return mqttClient; |
|||
} |
|||
|
|||
@Data |
|||
private class MqttMessageListener implements MqttHandler { |
|||
private final BlockingQueue<MqttEvent> events; |
|||
|
|||
private MqttMessageListener() { |
|||
events = new ArrayBlockingQueue<>(100); |
|||
} |
|||
|
|||
@Override |
|||
public void onMessage(String topic, ByteBuf message) { |
|||
log.info("MQTT message [{}], topic [{}]", message.toString(StandardCharsets.UTF_8), topic); |
|||
events.add(new MqttEvent(topic, message.toString(StandardCharsets.UTF_8))); |
|||
} |
|||
|
|||
} |
|||
|
|||
@Data |
|||
private class MqttEvent { |
|||
private final String topic; |
|||
private final String message; |
|||
} |
|||
|
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue