From 4b0554b274714657f70a03938b36e4f9551fce8a Mon Sep 17 00:00:00 2001 From: Viacheslav Kukhtyn Date: Tue, 16 Oct 2018 14:05:52 +0300 Subject: [PATCH 1/3] Integration-tests subproject --- msa/integration-tests/README.md | 18 +++ msa/integration-tests/pom.xml | 98 +++++++++++++++ .../server/msa/AbstractContainerTest.java | 112 ++++++++++++++++++ .../server/msa/ContainerTestSuite.java | 37 ++++++ .../org/thingsboard/server/msa/WsClient.java | 57 +++++++++ .../server/msa/WsTelemetryResponse.java | 40 +++++++ .../msa/connectivity/HttpClientTest.java | 57 +++++++++ .../msa/connectivity/MqttClientTest.java | 111 +++++++++++++++++ msa/pom.xml | 1 + 9 files changed, 531 insertions(+) create mode 100644 msa/integration-tests/README.md create mode 100644 msa/integration-tests/pom.xml create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java diff --git a/msa/integration-tests/README.md b/msa/integration-tests/README.md new file mode 100644 index 0000000000..5ae6354740 --- /dev/null +++ b/msa/integration-tests/README.md @@ -0,0 +1,18 @@ + +## Integration tests execution +To run the integration tests with using Docker, the local Docker images of Thingsboard's microservices should be built.
+- Build the local Docker images in the directory with the Thingsboard's main [pom.xml](./../../pom.xml): + + mvn clean install -Ddockerfile.skip=false +- Verify that the new local images were built: + + docker image ls +As result, in REPOSITORY column, next images should be present: + + local-maven-build/tb-node + local-maven-build/tb-web-ui + local-maven-build/tb-web-ui + +- Run the integration tests in the [msa/integration-tests](../integration-tests) directory: + + mvn clean install -Dintegrationtests.skip=false \ No newline at end of file diff --git a/msa/integration-tests/pom.xml b/msa/integration-tests/pom.xml new file mode 100644 index 0000000000..a1c24f39c6 --- /dev/null +++ b/msa/integration-tests/pom.xml @@ -0,0 +1,98 @@ + + + 4.0.0 + + + org.thingsboard + 2.2.0-SNAPSHOT + msa + + org.thingsboard.msa + integration-tests + + ThingsBoard Integration Tests + https://thingsboard.io + Project for ThingsBoard integration tests with using Docker + + + UTF-8 + ${basedir}/../.. + true + 1.9.1 + 1.3.9 + + + + + org.testcontainers + testcontainers + ${testcontainers.version} + + + org.java-websocket + Java-WebSocket + ${java-websocket.version} + + + io.takari.junit + takari-cpsuite + + + ch.qos.logback + logback-classic + + + com.google.code.gson + gson + + + org.apache.commons + commons-lang3 + + + com.google.guava + guava + + + org.thingsboard + netty-mqtt + + + org.thingsboard + tools + + + + + + + org.apache.maven.plugins + maven-surefire-plugin + + + **/*TestSuite.java + + ${integrationtests.skip} + + + + + + diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java new file mode 100644 index 0000000000..ab7f101f98 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java @@ -0,0 +1,112 @@ +/** + * Copyright © 2016-2018 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; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.ImmutableMap; +import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.RandomStringUtils; +import org.junit.*; +import org.thingsboard.client.tools.RestClient; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.id.DeviceId; + +import java.net.URI; +import java.net.URISyntaxException; +import java.util.List; +import java.util.Map; +import java.util.Random; +import java.util.concurrent.TimeUnit; + +@Slf4j +public abstract class AbstractContainerTest { + protected static String httpUrl; + protected static String wsUrl; + protected static RestClient restClient; + protected ObjectMapper mapper = new ObjectMapper(); + + @BeforeClass + public static void before() { + httpUrl = "http://localhost:" + ContainerTestSuite.composeContainer.getServicePort("tb-web-ui1", ContainerTestSuite.EXPOSED_PORT); + wsUrl = "ws://localhost:" + ContainerTestSuite.composeContainer.getServicePort("tb-web-ui1", ContainerTestSuite.EXPOSED_PORT); + restClient = new RestClient(httpUrl); + } + + protected Device createDevice(String name) { + return restClient.createDevice(name + RandomStringUtils.randomAlphanumeric(7), "DEFAULT"); + } + + protected WsClient subscribeToTelemetryWebSocket(DeviceId deviceId) throws URISyntaxException, InterruptedException { + WsClient mWs = new WsClient(new URI(wsUrl + "/api/ws/plugins/telemetry?token=" + restClient.getToken())); + mWs.connectBlocking(1, TimeUnit.SECONDS); + + JsonObject tsSubCmd = new JsonObject(); + tsSubCmd.addProperty("entityType", EntityType.DEVICE.name()); + tsSubCmd.addProperty("entityId", deviceId.toString()); + tsSubCmd.addProperty("scope", "LATEST_TELEMETRY"); + tsSubCmd.addProperty("cmdId", new Random().nextInt(100)); + tsSubCmd.addProperty("unsubscribe", false); + JsonArray wsTsSubCmds = new JsonArray(); + wsTsSubCmds.add(tsSubCmd); + JsonObject wsRequest = new JsonObject(); + wsRequest.add("tsSubCmds", wsTsSubCmds); + wsRequest.add("historyCmds", new JsonArray()); + wsRequest.add("attrSubCmds", new JsonArray()); + mWs.send(wsRequest.toString()); + return mWs; + } + + protected Map getExpectedLatestValues(long ts) { + return ImmutableMap.builder() + .put("booleanKey", ts) + .put("stringKey", ts) + .put("doubleKey", ts) + .put("longKey", ts) + .build(); + } + + protected boolean verify(WsTelemetryResponse wsTelemetryResponse, String key, Long expectedTs, String expectedValue) { + List list = wsTelemetryResponse.getDataValuesByKey(key); + return expectedTs.equals(list.get(0)) && expectedValue.equals(list.get(1)); + } + + protected boolean verify(WsTelemetryResponse wsTelemetryResponse, String key, String expectedValue) { + List list = wsTelemetryResponse.getDataValuesByKey(key); + return expectedValue.equals(list.get(1)); + } + + protected JsonObject createPayload(long ts) { + JsonObject values = createPayload(); + JsonObject payload = new JsonObject(); + payload.addProperty("ts", ts); + payload.add("values", values); + return payload; + } + + protected JsonObject createPayload() { + JsonObject values = new JsonObject(); + values.addProperty("stringKey", "value1"); + values.addProperty("booleanKey", true); + values.addProperty("doubleKey", 42.0); + values.addProperty("longKey", 73L); + + return values; + } + +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java new file mode 100644 index 0000000000..fd2de229d0 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2018 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; + +import org.junit.ClassRule; +import org.junit.extensions.cpsuite.ClasspathSuite; +import org.junit.runner.RunWith; +import org.testcontainers.containers.DockerComposeContainer; +import org.testcontainers.containers.wait.strategy.Wait; + +import java.io.File; + +@RunWith(ClasspathSuite.class) +@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.msa.*"}) +public class ContainerTestSuite { + static final int EXPOSED_PORT = 8080; + + @ClassRule + public static DockerComposeContainer composeContainer = new DockerComposeContainer(new File("./../docker/docker-compose.yml")) + .withPull(false) + .withLocalCompose(true) + .withTailChildContainers(true) + .withExposedService("tb-web-ui1", EXPOSED_PORT, Wait.forHttp("/login")); +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java new file mode 100644 index 0000000000..2f05eee4ce --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java @@ -0,0 +1,57 @@ +/** + * Copyright © 2016-2018 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; + +import org.java_websocket.client.WebSocketClient; +import org.java_websocket.handshake.ServerHandshake; + +import java.net.URI; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.BlockingQueue; + +public class WsClient extends WebSocketClient { + private final BlockingQueue events; + private String message; + + public WsClient(URI serverUri) { + super(serverUri); + events = new ArrayBlockingQueue<>(100); + } + + @Override + public void onOpen(ServerHandshake serverHandshake) { + } + + @Override + public void onMessage(String message) { + events.add(message); + this.message = message; + } + + @Override + public void onClose(int code, String reason, boolean remote) { + events.clear(); + } + + @Override + public void onError(Exception ex) { + ex.printStackTrace(); + } + + public String getLastMessage() { + return this.message; + } +} \ No newline at end of file diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java new file mode 100644 index 0000000000..834c5d1449 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2018 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; + +import lombok.Data; + +import java.io.Serializable; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +@Data +public class WsTelemetryResponse implements Serializable { + private int subscriptionId; + private int errorCode; + private String errorMsg; + private Map>> data; + private Map latestValues; + + public List getDataValuesByKey(String key) { + return data.entrySet().stream() + .filter(e -> e.getKey().equals(key)) + .flatMap(e -> e.getValue().stream().flatMap(Collection::stream)) + .collect(Collectors.toList()); + } +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java new file mode 100644 index 0000000000..7cc0a0f6e3 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java @@ -0,0 +1,57 @@ +/** + * Copyright © 2016-2018 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 org.junit.Assert; +import org.junit.Test; +import org.springframework.http.ResponseEntity; +import org.thingsboard.server.common.data.Device; +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.WsTelemetryResponse; + +import java.util.concurrent.TimeUnit; + +public class HttpClientTest extends AbstractContainerTest { + + @Test + public void telemetryUpdate() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + + Device device = createDevice("http_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); + ResponseEntity deviceTelemetryResponse = restClient.getRestTemplate() + .postForEntity(httpUrl + "/api/v1/{credentialsId}/telemetry", + mapper.readTree(createPayload().toString()), + ResponseEntity.class, + deviceCredentials.getCredentialsId()); + Assert.assertTrue(deviceTelemetryResponse.getStatusCode().is2xxSuccessful()); + TimeUnit.SECONDS.sleep(1); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + + Assert.assertEquals(getExpectedLatestValues(123456789L).keySet(), 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))); + + restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + } +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java new file mode 100644 index 0000000000..4ad638ed15 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java @@ -0,0 +1,111 @@ +/** + * Copyright © 2016-2018 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 io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import lombok.Data; +import org.junit.*; +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.security.DeviceCredentials; +import org.thingsboard.server.msa.AbstractContainerTest; +import org.thingsboard.server.msa.WsClient; +import org.thingsboard.server.msa.WsTelemetryResponse; + +import java.nio.charset.StandardCharsets; +import java.util.concurrent.*; + +public class MqttClientTest extends AbstractContainerTest { + + @Test + public void telemetryUpload() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); + MqttClient mqttClient = getMqttClient(deviceCredentials); + mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload().toString().getBytes())); + TimeUnit.SECONDS.sleep(1); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + + Assert.assertEquals(getExpectedLatestValues(123456789L).keySet(), 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))); + + restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + } + + @Test + public void telemetryUploadWithTs() throws Exception { + long ts = 1451649600512L; + + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); + MqttClient mqttClient = getMqttClient(deviceCredentials); + mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload(ts).toString().getBytes())); + TimeUnit.SECONDS.sleep(1); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + + 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))); + + restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + } + + private MqttClient getMqttClient(DeviceCredentials deviceCredentials) throws InterruptedException { + MqttMessageListener queue = new MqttMessageListener(); + MqttClientConfig clientConfig = new MqttClientConfig(); + clientConfig.setClientId("MQTT client from test"); + clientConfig.setUsername(deviceCredentials.getCredentialsId()); + MqttClient mqttClient = MqttClient.create(clientConfig, queue); + mqttClient.connect("localhost", 1883).sync(); + return mqttClient; + } + + @Data + private class MqttMessageListener implements MqttHandler { + private final BlockingQueue events; + + private MqttMessageListener() { + events = new ArrayBlockingQueue<>(100); + } + + @Override + public void onMessage(String topic, ByteBuf message) { + events.add(new MqttEvent(topic, message.toString(StandardCharsets.UTF_8))); + } + } + + @Data + private class MqttEvent { + private final String topic; + private final String message; + } +} diff --git a/msa/pom.xml b/msa/pom.xml index 21ddb4dafe..dd4c3653ba 100644 --- a/msa/pom.xml +++ b/msa/pom.xml @@ -40,6 +40,7 @@ js-executor web-ui tb-node + integration-tests From fb098e1bae920b89358596ee274f4fd960673952 Mon Sep 17 00:00:00 2001 From: Viacheslav Kukhtyn Date: Tue, 23 Oct 2018 16:41:43 +0300 Subject: [PATCH 2/3] Cover MQTT API with black box tests --- msa/integration-tests/README.md | 2 +- msa/integration-tests/pom.xml | 10 +- .../server/msa/AbstractContainerTest.java | 113 +++++-- .../server/msa/ContainerTestSuite.java | 3 +- .../org/thingsboard/server/msa/WsClient.java | 8 +- .../msa/connectivity/HttpClientTest.java | 14 +- .../msa/connectivity/MqttClientTest.java | 304 +++++++++++++++++- .../server/msa/mapper/AttributesResponse.java | 26 ++ .../msa/{ => mapper}/WsTelemetryResponse.java | 2 +- .../RpcResponseRuleChainMetadata.json | 59 ++++ 10 files changed, 484 insertions(+), 57 deletions(-) create mode 100644 msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java rename msa/integration-tests/src/test/java/org/thingsboard/server/msa/{ => mapper}/WsTelemetryResponse.java (96%) create mode 100644 msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json diff --git a/msa/integration-tests/README.md b/msa/integration-tests/README.md index 5ae6354740..93df2bf77f 100644 --- a/msa/integration-tests/README.md +++ b/msa/integration-tests/README.md @@ -15,4 +15,4 @@ As result, in REPOSITORY column, next images should be present: - Run the integration tests in the [msa/integration-tests](../integration-tests) directory: - mvn clean install -Dintegrationtests.skip=false \ No newline at end of file + mvn clean install -DintegrationTests.skip=false \ No newline at end of file diff --git a/msa/integration-tests/pom.xml b/msa/integration-tests/pom.xml index a1c24f39c6..a9c67ccf93 100644 --- a/msa/integration-tests/pom.xml +++ b/msa/integration-tests/pom.xml @@ -34,9 +34,10 @@ UTF-8 ${basedir}/../.. - true + true 1.9.1 1.3.9 + 4.5.6 @@ -50,6 +51,11 @@ Java-WebSocket ${java-websocket.version} + + org.apache.httpcomponents + httpclient + ${httpclient.version} + io.takari.junit takari-cpsuite @@ -89,7 +95,7 @@ **/*TestSuite.java - ${integrationtests.skip} + ${integrationTests.skip} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java index ab7f101f98..5ebab780be 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java @@ -21,55 +21,68 @@ import com.google.gson.JsonArray; import com.google.gson.JsonObject; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; +import org.apache.http.config.Registry; +import org.apache.http.config.RegistryBuilder; +import org.apache.http.conn.socket.ConnectionSocketFactory; +import org.apache.http.conn.ssl.SSLConnectionSocketFactory; +import org.apache.http.conn.ssl.TrustStrategy; +import org.apache.http.conn.ssl.X509HostnameVerifier; +import org.apache.http.impl.client.CloseableHttpClient; +import org.apache.http.impl.client.HttpClients; +import org.apache.http.impl.conn.PoolingHttpClientConnectionManager; +import org.apache.http.ssl.SSLContextBuilder; +import org.apache.http.ssl.SSLContexts; import org.junit.*; +import org.springframework.http.client.HttpComponentsClientHttpRequestFactory; import org.thingsboard.client.tools.RestClient; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.msa.mapper.WsTelemetryResponse; +import javax.net.ssl.*; import java.net.URI; -import java.net.URISyntaxException; +import java.security.cert.X509Certificate; import java.util.List; import java.util.Map; import java.util.Random; -import java.util.concurrent.TimeUnit; @Slf4j public abstract class AbstractContainerTest { - protected static String httpUrl; - protected static String wsUrl; + protected static final String HTTPS_URL = "https://localhost"; + protected static final String WSS_URL = "wss://localhost"; protected static RestClient restClient; protected ObjectMapper mapper = new ObjectMapper(); @BeforeClass - public static void before() { - httpUrl = "http://localhost:" + ContainerTestSuite.composeContainer.getServicePort("tb-web-ui1", ContainerTestSuite.EXPOSED_PORT); - wsUrl = "ws://localhost:" + ContainerTestSuite.composeContainer.getServicePort("tb-web-ui1", ContainerTestSuite.EXPOSED_PORT); - restClient = new RestClient(httpUrl); + public static void before() throws Exception { + restClient = new RestClient(HTTPS_URL); + restClient.getRestTemplate().setRequestFactory(getRequestFactoryForSelfSignedCert()); } protected Device createDevice(String name) { return restClient.createDevice(name + RandomStringUtils.randomAlphanumeric(7), "DEFAULT"); } - protected WsClient subscribeToTelemetryWebSocket(DeviceId deviceId) throws URISyntaxException, InterruptedException { - WsClient mWs = new WsClient(new URI(wsUrl + "/api/ws/plugins/telemetry?token=" + restClient.getToken())); - mWs.connectBlocking(1, TimeUnit.SECONDS); - - JsonObject tsSubCmd = new JsonObject(); - tsSubCmd.addProperty("entityType", EntityType.DEVICE.name()); - tsSubCmd.addProperty("entityId", deviceId.toString()); - tsSubCmd.addProperty("scope", "LATEST_TELEMETRY"); - tsSubCmd.addProperty("cmdId", new Random().nextInt(100)); - tsSubCmd.addProperty("unsubscribe", false); - JsonArray wsTsSubCmds = new JsonArray(); - wsTsSubCmds.add(tsSubCmd); + protected WsClient subscribeToWebSocket(DeviceId deviceId, String scope, CmdsType property) throws Exception { + WsClient wsClient = new WsClient(new URI(WSS_URL + "/api/ws/plugins/telemetry?token=" + restClient.getToken())); + SSLContextBuilder builder = SSLContexts.custom(); + builder.loadTrustMaterial(null, (TrustStrategy) (chain, authType) -> true); + wsClient.setSocket(builder.build().getSocketFactory().createSocket()); + wsClient.connectBlocking(); + + JsonObject cmdsObject = new JsonObject(); + cmdsObject.addProperty("entityType", EntityType.DEVICE.name()); + cmdsObject.addProperty("entityId", deviceId.toString()); + cmdsObject.addProperty("scope", scope); + cmdsObject.addProperty("cmdId", new Random().nextInt(100)); + + JsonArray cmd = new JsonArray(); + cmd.add(cmdsObject); JsonObject wsRequest = new JsonObject(); - wsRequest.add("tsSubCmds", wsTsSubCmds); - wsRequest.add("historyCmds", new JsonArray()); - wsRequest.add("attrSubCmds", new JsonArray()); - mWs.send(wsRequest.toString()); - return mWs; + wsRequest.add(property.toString(), cmd); + wsClient.send(wsRequest.toString()); + return wsClient; } protected Map getExpectedLatestValues(long ts) { @@ -109,4 +122,54 @@ public abstract class AbstractContainerTest { return values; } + protected enum CmdsType { + TS_SUB_CMDS("tsSubCmds"), + HISTORY_CMDS("historyCmds"), + ATTR_SUB_CMDS("attrSubCmds"); + + private final String text; + + CmdsType(final String text) { + this.text = text; + } + + @Override + public String toString() { + return text; + } + } + + private static HttpComponentsClientHttpRequestFactory getRequestFactoryForSelfSignedCert() throws Exception { + SSLContextBuilder builder = SSLContexts.custom(); + builder.loadTrustMaterial(null, (TrustStrategy) (chain, authType) -> true); + SSLContext sslContext = builder.build(); + SSLConnectionSocketFactory sslSelfSigned = new SSLConnectionSocketFactory(sslContext, new X509HostnameVerifier() { + @Override + public void verify(String host, SSLSocket ssl) { + } + + @Override + public void verify(String host, X509Certificate cert) { + } + + @Override + public void verify(String host, String[] cns, String[] subjectAlts) { + } + + @Override + public boolean verify(String s, SSLSession sslSession) { + return true; + } + }); + + Registry socketFactoryRegistry = RegistryBuilder + .create() + .register("https", sslSelfSigned) + .build(); + + PoolingHttpClientConnectionManager cm = new PoolingHttpClientConnectionManager(socketFactoryRegistry); + CloseableHttpClient httpClient = HttpClients.custom().setConnectionManager(cm).build(); + return new HttpComponentsClientHttpRequestFactory(httpClient); + } + } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java index fd2de229d0..0629960b24 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -26,12 +26,11 @@ import java.io.File; @RunWith(ClasspathSuite.class) @ClasspathSuite.ClassnameFilters({"org.thingsboard.server.msa.*"}) public class ContainerTestSuite { - static final int EXPOSED_PORT = 8080; @ClassRule public static DockerComposeContainer composeContainer = new DockerComposeContainer(new File("./../docker/docker-compose.yml")) .withPull(false) .withLocalCompose(true) .withTailChildContainers(true) - .withExposedService("tb-web-ui1", EXPOSED_PORT, Wait.forHttp("/login")); + .withExposedService("tb-web-ui1", 8080, Wait.forHttp("/login")); } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java index 2f05eee4ce..5ef238f8aa 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java @@ -19,16 +19,12 @@ import org.java_websocket.client.WebSocketClient; import org.java_websocket.handshake.ServerHandshake; import java.net.URI; -import java.util.concurrent.ArrayBlockingQueue; -import java.util.concurrent.BlockingQueue; public class WsClient extends WebSocketClient { - private final BlockingQueue events; private String message; public WsClient(URI serverUri) { super(serverUri); - events = new ArrayBlockingQueue<>(100); } @Override @@ -37,13 +33,11 @@ public class WsClient extends WebSocketClient { @Override public void onMessage(String message) { - events.add(message); this.message = message; } @Override public void onClose(int code, String reason, boolean remote) { - events.clear(); } @Override @@ -54,4 +48,4 @@ public class WsClient extends WebSocketClient { public String getLastMessage() { return this.message; } -} \ No newline at end of file +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java index 7cc0a0f6e3..f8e1041252 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.msa.connectivity; +import com.google.common.collect.Sets; import org.junit.Assert; import org.junit.Test; import org.springframework.http.ResponseEntity; @@ -22,7 +23,7 @@ import org.thingsboard.server.common.data.Device; 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.WsTelemetryResponse; +import org.thingsboard.server.msa.mapper.WsTelemetryResponse; import java.util.concurrent.TimeUnit; @@ -35,23 +36,24 @@ public class HttpClientTest extends AbstractContainerTest { Device device = createDevice("http_"); DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); - WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); + WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); ResponseEntity deviceTelemetryResponse = restClient.getRestTemplate() - .postForEntity(httpUrl + "/api/v1/{credentialsId}/telemetry", + .postForEntity(HTTPS_URL + "/api/v1/{credentialsId}/telemetry", mapper.readTree(createPayload().toString()), ResponseEntity.class, deviceCredentials.getCredentialsId()); Assert.assertTrue(deviceTelemetryResponse.getStatusCode().is2xxSuccessful()); TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); - Assert.assertEquals(getExpectedLatestValues(123456789L).keySet(), actualLatestTelemetry.getLatestValues().keySet()); + 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))); - restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); } } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java index 4ad638ed15..eae98b9f1e 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java @@ -15,20 +15,39 @@ */ 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 io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import lombok.Data; +import org.apache.commons.lang3.RandomStringUtils; import org.junit.*; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.http.HttpMethod; +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.RuleChainId; +import org.thingsboard.server.common.data.page.TextPageData; +import org.thingsboard.server.common.data.rule.NodeConnectionInfo; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleNode; 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.WsTelemetryResponse; +import org.thingsboard.server.msa.mapper.AttributesResponse; +import org.thingsboard.server.msa.mapper.WsTelemetryResponse; +import java.io.IOException; import java.nio.charset.StandardCharsets; +import java.util.*; import java.util.concurrent.*; public class MqttClientTest extends AbstractContainerTest { @@ -39,20 +58,22 @@ public class MqttClientTest extends AbstractContainerTest { Device device = createDevice("mqtt_"); DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); - WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); - MqttClient mqttClient = getMqttClient(deviceCredentials); + WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); + MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload().toString().getBytes())); TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); - Assert.assertEquals(getExpectedLatestValues(123456789L).keySet(), actualLatestTelemetry.getLatestValues().keySet()); + 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))); - restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); } @Test @@ -63,12 +84,13 @@ public class MqttClientTest extends AbstractContainerTest { Device device = createDevice("mqtt_"); DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); - WsClient mWs = subscribeToTelemetryWebSocket(device.getId()); - MqttClient mqttClient = getMqttClient(deviceCredentials); + WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); + MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload(ts).toString().getBytes())); TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(mWs.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + Assert.assertEquals(4, actualLatestTelemetry.getData().size()); Assert.assertEquals(getExpectedLatestValues(ts), actualLatestTelemetry.getLatestValues()); Assert.assertTrue(verify(actualLatestTelemetry, "booleanKey", ts, Boolean.TRUE.toString())); @@ -76,15 +98,271 @@ public class MqttClientTest extends AbstractContainerTest { Assert.assertTrue(verify(actualLatestTelemetry, "doubleKey", ts, Double.toString(42.0))); Assert.assertTrue(verify(actualLatestTelemetry, "longKey", ts, Long.toString(73))); - restClient.getRestTemplate().delete(httpUrl + "/api/device/" + device.getId()); + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); } - private MqttClient getMqttClient(DeviceCredentials deviceCredentials) throws InterruptedException { - MqttMessageListener queue = new MqttMessageListener(); + @Test + public void publishAttributeUpdateToServer() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + WsClient wsClient = subscribeToWebSocket(device.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); + MqttMessageListener listener = new MqttMessageListener(); + MqttClient mqttClient = getMqttClient(deviceCredentials, listener); + JsonObject clientAttributes = new JsonObject(); + clientAttributes.addProperty("attr1", "value1"); + clientAttributes.addProperty("attr2", true); + clientAttributes.addProperty("attr3", 42.0); + clientAttributes.addProperty("attr4", 73); + mqttClient.publish("v1/devices/me/attributes", Unpooled.wrappedBuffer(clientAttributes.toString().getBytes())); + TimeUnit.SECONDS.sleep(1); + WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + + 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))); + + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); + } + + @Test + public void requestAttributeValuesFromServer() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + MqttMessageListener listener = new MqttMessageListener(); + MqttClient mqttClient = getMqttClient(deviceCredentials, listener); + + // Add a new client attribute + JsonObject clientAttributes = new JsonObject(); + String clientAttributeValue = RandomStringUtils.randomAlphanumeric(8); + clientAttributes.addProperty("clientAttr", clientAttributeValue); + mqttClient.publish("v1/devices/me/attributes", Unpooled.wrappedBuffer(clientAttributes.toString().getBytes())); + + // 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, + device.getId()); + Assert.assertTrue(sharedAttributesResponse.getStatusCode().is2xxSuccessful()); + + // Subscribe to attributes response + mqttClient.on("v1/devices/me/attributes/response/+", listener); + // Request attributes + JsonObject request = new JsonObject(); + request.addProperty("clientKeys", "clientAttr"); + request.addProperty("sharedKeys", "sharedAttr"); + mqttClient.publish("v1/devices/me/attributes/request/" + new Random().nextInt(100), Unpooled.wrappedBuffer(request.toString().getBytes())); + MqttEvent event = listener.getEvents().poll(10, TimeUnit.SECONDS); + AttributesResponse attributes = mapper.readValue(Objects.requireNonNull(event).getMessage(), AttributesResponse.class); + + Assert.assertEquals(1, attributes.getClient().size()); + Assert.assertEquals(clientAttributeValue, attributes.getClient().get("clientAttr")); + + Assert.assertEquals(1, attributes.getShared().size()); + Assert.assertEquals(sharedAttributeValue, attributes.getShared().get("sharedAttr")); + + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); + } + + @Test + public void subscribeToAttributeUpdatesFromServer() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + MqttMessageListener listener = new MqttMessageListener(); + MqttClient mqttClient = getMqttClient(deviceCredentials, listener); + mqttClient.on("v1/devices/me/attributes", listener); + + String sharedAttributeName = "sharedAttr"; + + // Add a new shared attribute + JsonObject sharedAttributes = new JsonObject(); + String sharedAttributeValue = RandomStringUtils.randomAlphanumeric(8); + sharedAttributes.addProperty(sharedAttributeName, sharedAttributeValue); + ResponseEntity sharedAttributesResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", + mapper.readTree(sharedAttributes.toString()), ResponseEntity.class, + device.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(sharedAttributeName).asText()); + + // Update the shared attribute value + JsonObject updatedSharedAttributes = new JsonObject(); + String updatedSharedAttributeValue = RandomStringUtils.randomAlphanumeric(8); + updatedSharedAttributes.addProperty(sharedAttributeName, updatedSharedAttributeValue); + ResponseEntity updatedSharedAttributesResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", + mapper.readTree(updatedSharedAttributes.toString()), ResponseEntity.class, + device.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(sharedAttributeName).asText()); + + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); + } + + @Test + public void serverSideRpc() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + MqttMessageListener listener = new MqttMessageListener(); + MqttClient mqttClient = getMqttClient(deviceCredentials, listener); + mqttClient.on("v1/devices/me/rpc/request/+", listener); + + // Send an RPC from the server + JsonObject serverRpcPayload = new JsonObject(); + serverRpcPayload.addProperty("method", "getValue"); + serverRpcPayload.addProperty("params", true); + serverRpcPayload.addProperty("timeout", 1000); + ListeningExecutorService service = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); + ListenableFuture future = service.submit(() -> { + try { + return restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/plugins/rpc/twoway/{deviceId}", + mapper.readTree(serverRpcPayload.toString()), String.class, + device.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.assertEquals("{\"method\":\"getValue\",\"params\":true}", Objects.requireNonNull(requestFromServer).getMessage()); + + Integer requestId = Integer.valueOf(Objects.requireNonNull(requestFromServer).getTopic().substring("v1/devices/me/rpc/request/".length())); + JsonObject clientResponse = new JsonObject(); + clientResponse.addProperty("response", "someResponse"); + // Send a response to the server's RPC request + mqttClient.publish("v1/devices/me/rpc/response/" + requestId, Unpooled.wrappedBuffer(clientResponse.toString().getBytes())); + + ResponseEntity serverResponse = future.get(5, TimeUnit.SECONDS); + Assert.assertTrue(serverResponse.getStatusCode().is2xxSuccessful()); + Assert.assertEquals(clientResponse.toString(), serverResponse.getBody()); + + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); + } + + @Test + public void clientSideRpc() throws Exception { + restClient.login("tenant@thingsboard.org", "tenant"); + Device device = createDevice("mqtt_"); + DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId()); + + MqttMessageListener listener = new MqttMessageListener(); + MqttClient mqttClient = getMqttClient(deviceCredentials, listener); + mqttClient.on("v1/devices/me/rpc/request/+", listener); + + // Get the default rule chain id to make it root again after test finished + RuleChainId defaultRuleChainId = getDefaultRuleChainId(); + + // Create a new root rule chain + RuleChainId ruleChainId = createRootRuleChainForRpcResponse(); + + // Send the request to the server + JsonObject clientRequest = new JsonObject(); + clientRequest.addProperty("method", "getResponse"); + clientRequest.addProperty("params", true); + Integer requestId = 42; + mqttClient.publish("v1/devices/me/rpc/request/" + requestId, Unpooled.wrappedBuffer(clientRequest.toString().getBytes())); + + // Check the response from the server + TimeUnit.SECONDS.sleep(1); + MqttEvent responseFromServer = listener.getEvents().poll(1, TimeUnit.SECONDS); + Integer responseId = Integer.valueOf(Objects.requireNonNull(responseFromServer).getTopic().substring("v1/devices/me/rpc/response/".length())); + Assert.assertEquals(requestId, responseId); + Assert.assertEquals("requestReceived", mapper.readTree(responseFromServer.getMessage()).get("response").asText()); + + // Make the default rule chain a root again + ResponseEntity rootRuleChainResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/ruleChain/{ruleChainId}/root", + null, + RuleChain.class, + defaultRuleChainId); + Assert.assertTrue(rootRuleChainResponse.getStatusCode().is2xxSuccessful()); + + // Delete the created rule chain + restClient.getRestTemplate().delete(HTTPS_URL + "/api/ruleChain/{ruleChainId}", ruleChainId); + restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId()); + } + + private RuleChainId createRootRuleChainForRpcResponse() throws Exception { + RuleChain newRuleChain = new RuleChain(); + newRuleChain.setName("testRuleChain"); + ResponseEntity ruleChainResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/ruleChain", + newRuleChain, + RuleChain.class); + Assert.assertTrue(ruleChainResponse.getStatusCode().is2xxSuccessful()); + RuleChain ruleChain = ruleChainResponse.getBody(); + + JsonNode configuration = mapper.readTree(this.getClass().getClassLoader().getResourceAsStream("RpcResponseRuleChainMetadata.json")); + RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); + ruleChainMetaData.setRuleChainId(ruleChain.getId()); + ruleChainMetaData.setFirstNodeIndex(configuration.get("firstNodeIndex").asInt()); + ruleChainMetaData.setNodes(Arrays.asList(mapper.treeToValue(configuration.get("nodes"), RuleNode[].class))); + ruleChainMetaData.setConnections(Arrays.asList(mapper.treeToValue(configuration.get("connections"), NodeConnectionInfo[].class))); + + ResponseEntity ruleChainMetadataResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/ruleChain/metadata", + ruleChainMetaData, + RuleChainMetaData.class); + Assert.assertTrue(ruleChainMetadataResponse.getStatusCode().is2xxSuccessful()); + + // Set a new rule chain as root + ResponseEntity rootRuleChainResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/ruleChain/{ruleChainId}/root", + null, + RuleChain.class, + ruleChain.getId()); + Assert.assertTrue(rootRuleChainResponse.getStatusCode().is2xxSuccessful()); + + return ruleChain.getId(); + } + + private RuleChainId getDefaultRuleChainId() { + ResponseEntity> ruleChains = restClient.getRestTemplate().exchange( + HTTPS_URL + "/api/ruleChains?limit=40&textSearch=", + HttpMethod.GET, + null, + new ParameterizedTypeReference>() { + }); + + Optional defaultRuleChain = ruleChains.getBody().getData() + .stream() + .filter(RuleChain::isRoot) + .findFirst(); + if (!defaultRuleChain.isPresent()) { + Assert.fail("Root rule chain wasn't found"); + } + return defaultRuleChain.get().getId(); + } + + private MqttClient getMqttClient(DeviceCredentials deviceCredentials, MqttMessageListener listener) throws InterruptedException { MqttClientConfig clientConfig = new MqttClientConfig(); clientConfig.setClientId("MQTT client from test"); clientConfig.setUsername(deviceCredentials.getCredentialsId()); - MqttClient mqttClient = MqttClient.create(clientConfig, queue); + MqttClient mqttClient = MqttClient.create(clientConfig, listener); mqttClient.connect("localhost", 1883).sync(); return mqttClient; } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java new file mode 100644 index 0000000000..f9774eef20 --- /dev/null +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2018 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.mapper; + +import lombok.Data; + +import java.util.Map; + +@Data +public class AttributesResponse { + private Map client; + private Map shared; +} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java similarity index 96% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java rename to msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java index 834c5d1449..b22244f3a7 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsTelemetryResponse.java +++ b/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.msa; +package org.thingsboard.server.msa.mapper; import lombok.Data; diff --git a/msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json b/msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json new file mode 100644 index 0000000000..09178ef781 --- /dev/null +++ b/msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json @@ -0,0 +1,59 @@ +{ + "firstNodeIndex": 0, + "nodes": [ + { + "additionalInfo": { + "layoutX": 325, + "layoutY": 150 + }, + "type": "org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode", + "name": "msgTypeSwitch", + "debugMode": true, + "configuration": { + "version": 0 + } + }, + { + "additionalInfo": { + "layoutX": 60, + "layoutY": 300 + }, + "type": "org.thingsboard.rule.engine.transform.TbTransformMsgNode", + "name": "formResponse", + "debugMode": true, + "configuration": { + "jsScript": "if (msg.method == \"getResponse\") {\n return {msg: {\"response\": \"requestReceived\"}, metadata: metadata, msgType: msgType};\n}\n\nreturn {msg: msg, metadata: metadata, msgType: msgType};" + } + }, + { + "additionalInfo": { + "layoutX": 450, + "layoutY": 300 + }, + "type": "org.thingsboard.rule.engine.rpc.TbSendRPCReplyNode", + "name": "rpcReply", + "debugMode": true, + "configuration": { + "requestIdMetaDataAttribute": "requestId" + } + } + ], + "connections": [ + { + "fromIndex": 0, + "toIndex": 1, + "type": "RPC Request from Device" + }, + { + "fromIndex": 1, + "toIndex": 2, + "type": "Success" + }, + { + "fromIndex": 1, + "toIndex": 2, + "type": "Failure" + } + ], + "ruleChainConnections": null +} \ No newline at end of file From 32dd7f0437053fdbb42227ff939dc9f2e7be17d1 Mon Sep 17 00:00:00 2001 From: Viacheslav Kukhtyn Date: Wed, 31 Oct 2018 13:07:43 +0200 Subject: [PATCH 3/3] Black box tests update --- msa/black-box-tests/README.md | 21 ++++++++++++ .../pom.xml | 10 +++--- .../server/msa/AbstractContainerTest.java | 0 .../server/msa/ContainerTestSuite.java | 9 +++-- .../org/thingsboard/server/msa/WsClient.java | 33 ++++++++++++++++--- .../msa/connectivity/HttpClientTest.java | 8 ++--- .../msa/connectivity/MqttClientTest.java | 25 +++++++------- .../server/msa/mapper/AttributesResponse.java | 0 .../msa/mapper/WsTelemetryResponse.java | 0 .../RpcResponseRuleChainMetadata.json | 0 msa/integration-tests/README.md | 18 ---------- msa/pom.xml | 4 +-- 12 files changed, 80 insertions(+), 48 deletions(-) create mode 100644 msa/black-box-tests/README.md rename msa/{integration-tests => black-box-tests}/pom.xml (92%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java (100%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java (80%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/WsClient.java (52%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java (90%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java (95%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java (100%) rename msa/{integration-tests => black-box-tests}/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java (100%) rename msa/{integration-tests => black-box-tests}/src/test/resources/RpcResponseRuleChainMetadata.json (100%) delete mode 100644 msa/integration-tests/README.md diff --git a/msa/black-box-tests/README.md b/msa/black-box-tests/README.md new file mode 100644 index 0000000000..c26d9c5fc6 --- /dev/null +++ b/msa/black-box-tests/README.md @@ -0,0 +1,21 @@ + +## Black box tests execution +To run the black box tests with using Docker, the local Docker images of Thingsboard's microservices should be built.
+- Build the local Docker images in the directory with the Thingsboard's main [pom.xml](./../../pom.xml): + + mvn clean install -Ddockerfile.skip=false +- Verify that the new local images were built: + + docker image ls +As result, in REPOSITORY column, next images should be present: + + thingsboard/tb-coap-transport + thingsboard/tb-http-transport + thingsboard/tb-mqtt-transport + thingsboard/tb-node + thingsboard/tb-web-ui + thingsboard/tb-js-executor + +- Run the black box tests in the [msa/black-box-tests](../black-box-tests) directory: + + mvn clean install -DblackBoxTests.skip=false \ No newline at end of file diff --git a/msa/integration-tests/pom.xml b/msa/black-box-tests/pom.xml similarity index 92% rename from msa/integration-tests/pom.xml rename to msa/black-box-tests/pom.xml index a9c67ccf93..2af4d2c14e 100644 --- a/msa/integration-tests/pom.xml +++ b/msa/black-box-tests/pom.xml @@ -25,16 +25,16 @@ msa org.thingsboard.msa - integration-tests + black-box-tests - ThingsBoard Integration Tests + ThingsBoard Black Box Tests https://thingsboard.io - Project for ThingsBoard integration tests with using Docker + Project for ThingsBoard black box testing with using Docker UTF-8 ${basedir}/../.. - true + true 1.9.1 1.3.9 4.5.6 @@ -95,7 +95,7 @@ **/*TestSuite.java - ${integrationTests.skip} + ${blackBoxTests.skip} diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java similarity index 100% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java similarity index 80% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java index 0629960b24..495fd94d2e 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -22,15 +22,18 @@ import org.testcontainers.containers.DockerComposeContainer; import org.testcontainers.containers.wait.strategy.Wait; import java.io.File; +import java.time.Duration; @RunWith(ClasspathSuite.class) -@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.msa.*"}) +@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.msa.*Test"}) public class ContainerTestSuite { @ClassRule - public static DockerComposeContainer composeContainer = new DockerComposeContainer(new File("./../docker/docker-compose.yml")) + public static DockerComposeContainer composeContainer = new DockerComposeContainer( + new File("./../../docker/docker-compose.yml"), + new File("./../../docker/docker-compose.postgres.yml")) .withPull(false) .withLocalCompose(true) .withTailChildContainers(true) - .withExposedService("tb-web-ui1", 8080, Wait.forHttp("/login")); + .withExposedService("tb-web-ui1", 8080, Wait.forHttp("/login").withStartupTimeout(Duration.ofSeconds(120))); } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java similarity index 52% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java index 5ef238f8aa..a9835ed618 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/WsClient.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java @@ -15,13 +15,23 @@ */ package org.thingsboard.server.msa; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; import org.java_websocket.client.WebSocketClient; import org.java_websocket.handshake.ServerHandshake; +import org.thingsboard.server.msa.mapper.WsTelemetryResponse; +import java.io.IOException; import java.net.URI; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +@Slf4j public class WsClient extends WebSocketClient { - private String message; + private static final ObjectMapper mapper = new ObjectMapper(); + private WsTelemetryResponse message; + + private CountDownLatch latch = new CountDownLatch(1);; public WsClient(URI serverUri) { super(serverUri); @@ -33,11 +43,20 @@ public class WsClient extends WebSocketClient { @Override public void onMessage(String message) { - this.message = message; + try { + WsTelemetryResponse response = mapper.readValue(message, WsTelemetryResponse.class); + if (!response.getData().isEmpty()) { + this.message = response; + latch.countDown(); + } + } catch (IOException e) { + log.error("ws message can't be read"); + } } @Override public void onClose(int code, String reason, boolean remote) { + log.info("ws is closed, due to [{}]", reason); } @Override @@ -45,7 +64,13 @@ public class WsClient extends WebSocketClient { ex.printStackTrace(); } - public String getLastMessage() { - return this.message; + public WsTelemetryResponse getLastMessage() { + try { + latch.await(10, TimeUnit.SECONDS); + return this.message; + } catch (InterruptedException e) { + log.error("Timeout, ws message wasn't received"); + } + return null; } } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java similarity index 90% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java index f8e1041252..a6e89de185 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java @@ -25,12 +25,10 @@ import org.thingsboard.server.msa.AbstractContainerTest; import org.thingsboard.server.msa.WsClient; import org.thingsboard.server.msa.mapper.WsTelemetryResponse; -import java.util.concurrent.TimeUnit; - public class HttpClientTest extends AbstractContainerTest { @Test - public void telemetryUpdate() throws Exception { + public void telemetryUpload() throws Exception { restClient.login("tenant@thingsboard.org", "tenant"); Device device = createDevice("http_"); @@ -43,8 +41,8 @@ public class HttpClientTest extends AbstractContainerTest { ResponseEntity.class, deviceCredentials.getCredentialsId()); Assert.assertTrue(deviceTelemetryResponse.getStatusCode().is2xxSuccessful()); - TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + wsClient.closeBlocking(); Assert.assertEquals(Sets.newHashSet("booleanKey", "stringKey", "doubleKey", "longKey"), actualLatestTelemetry.getLatestValues().keySet()); diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java similarity index 95% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java index eae98b9f1e..d889d2c52e 100644 --- a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java @@ -23,7 +23,9 @@ import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.JsonObject; 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.*; import org.springframework.core.ParameterizedTypeReference; @@ -50,6 +52,7 @@ import java.nio.charset.StandardCharsets; import java.util.*; import java.util.concurrent.*; +@Slf4j public class MqttClientTest extends AbstractContainerTest { @Test @@ -61,8 +64,8 @@ public class MqttClientTest extends AbstractContainerTest { WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload().toString().getBytes())); - TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); Assert.assertEquals(Sets.newHashSet("booleanKey", "stringKey", "doubleKey", "longKey"), @@ -87,8 +90,8 @@ public class MqttClientTest extends AbstractContainerTest { WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload(ts).toString().getBytes())); - TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); Assert.assertEquals(getExpectedLatestValues(ts), actualLatestTelemetry.getLatestValues()); @@ -116,8 +119,8 @@ public class MqttClientTest extends AbstractContainerTest { clientAttributes.addProperty("attr3", 42.0); clientAttributes.addProperty("attr4", 73); mqttClient.publish("v1/devices/me/attributes", Unpooled.wrappedBuffer(clientAttributes.toString().getBytes())); - TimeUnit.SECONDS.sleep(1); - WsTelemetryResponse actualLatestTelemetry = mapper.readValue(wsClient.getLastMessage(), WsTelemetryResponse.class); + WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); Assert.assertEquals(Sets.newHashSet("attr1", "attr2", "attr3", "attr4"), @@ -157,7 +160,7 @@ public class MqttClientTest extends AbstractContainerTest { Assert.assertTrue(sharedAttributesResponse.getStatusCode().is2xxSuccessful()); // Subscribe to attributes response - mqttClient.on("v1/devices/me/attributes/response/+", listener); + mqttClient.on("v1/devices/me/attributes/response/+", listener, MqttQoS.AT_LEAST_ONCE); // Request attributes JsonObject request = new JsonObject(); request.addProperty("clientKeys", "clientAttr"); @@ -183,7 +186,7 @@ public class MqttClientTest extends AbstractContainerTest { MqttMessageListener listener = new MqttMessageListener(); MqttClient mqttClient = getMqttClient(deviceCredentials, listener); - mqttClient.on("v1/devices/me/attributes", listener); + mqttClient.on("v1/devices/me/attributes", listener, MqttQoS.AT_LEAST_ONCE); String sharedAttributeName = "sharedAttr"; @@ -226,13 +229,12 @@ public class MqttClientTest extends AbstractContainerTest { MqttMessageListener listener = new MqttMessageListener(); MqttClient mqttClient = getMqttClient(deviceCredentials, listener); - mqttClient.on("v1/devices/me/rpc/request/+", listener); + mqttClient.on("v1/devices/me/rpc/request/+", listener, MqttQoS.AT_LEAST_ONCE); // Send an RPC from the server JsonObject serverRpcPayload = new JsonObject(); serverRpcPayload.addProperty("method", "getValue"); serverRpcPayload.addProperty("params", true); - serverRpcPayload.addProperty("timeout", 1000); ListeningExecutorService service = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); ListenableFuture future = service.submit(() -> { try { @@ -271,7 +273,7 @@ public class MqttClientTest extends AbstractContainerTest { MqttMessageListener listener = new MqttMessageListener(); MqttClient mqttClient = getMqttClient(deviceCredentials, listener); - mqttClient.on("v1/devices/me/rpc/request/+", listener); + mqttClient.on("v1/devices/me/rpc/request/+", listener, MqttQoS.AT_LEAST_ONCE); // Get the default rule chain id to make it root again after test finished RuleChainId defaultRuleChainId = getDefaultRuleChainId(); @@ -377,6 +379,7 @@ public class MqttClientTest extends AbstractContainerTest { @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))); } } diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java similarity index 100% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/mapper/AttributesResponse.java diff --git a/msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java similarity index 100% rename from msa/integration-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java rename to msa/black-box-tests/src/test/java/org/thingsboard/server/msa/mapper/WsTelemetryResponse.java diff --git a/msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json b/msa/black-box-tests/src/test/resources/RpcResponseRuleChainMetadata.json similarity index 100% rename from msa/integration-tests/src/test/resources/RpcResponseRuleChainMetadata.json rename to msa/black-box-tests/src/test/resources/RpcResponseRuleChainMetadata.json diff --git a/msa/integration-tests/README.md b/msa/integration-tests/README.md deleted file mode 100644 index 93df2bf77f..0000000000 --- a/msa/integration-tests/README.md +++ /dev/null @@ -1,18 +0,0 @@ - -## Integration tests execution -To run the integration tests with using Docker, the local Docker images of Thingsboard's microservices should be built.
-- Build the local Docker images in the directory with the Thingsboard's main [pom.xml](./../../pom.xml): - - mvn clean install -Ddockerfile.skip=false -- Verify that the new local images were built: - - docker image ls -As result, in REPOSITORY column, next images should be present: - - local-maven-build/tb-node - local-maven-build/tb-web-ui - local-maven-build/tb-web-ui - -- Run the integration tests in the [msa/integration-tests](../integration-tests) directory: - - mvn clean install -DintegrationTests.skip=false \ No newline at end of file diff --git a/msa/pom.xml b/msa/pom.xml index 41926608be..5a9d9abd80 100644 --- a/msa/pom.xml +++ b/msa/pom.xml @@ -16,7 +16,7 @@ --> + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 org.thingsboard @@ -41,7 +41,7 @@ web-ui tb-node transport - integration-tests + black-box-tests