|
|
|
@ -16,7 +16,6 @@ |
|
|
|
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; |
|
|
|
@ -28,15 +27,17 @@ import io.netty.buffer.Unpooled; |
|
|
|
import io.netty.handler.codec.mqtt.MqttQoS; |
|
|
|
import lombok.Data; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.junit.After; |
|
|
|
import org.junit.Assert; |
|
|
|
import org.junit.Before; |
|
|
|
import org.junit.Test; |
|
|
|
import org.springframework.http.ResponseEntity; |
|
|
|
import org.springframework.http.HttpStatus; |
|
|
|
import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; |
|
|
|
import org.testng.annotations.AfterMethod; |
|
|
|
import org.testng.annotations.BeforeMethod; |
|
|
|
import org.testng.annotations.Test; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
|
|
import org.thingsboard.mqtt.MqttClient; |
|
|
|
import org.thingsboard.mqtt.MqttClientConfig; |
|
|
|
import org.thingsboard.mqtt.MqttHandler; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
import org.thingsboard.server.common.data.Device; |
|
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
@ -50,9 +51,9 @@ import org.thingsboard.server.msa.mapper.WsTelemetryResponse; |
|
|
|
|
|
|
|
import java.io.IOException; |
|
|
|
import java.nio.charset.StandardCharsets; |
|
|
|
import java.util.Arrays; |
|
|
|
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; |
|
|
|
@ -60,28 +61,34 @@ import java.util.concurrent.ExecutionException; |
|
|
|
import java.util.concurrent.Executors; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
|
|
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
|
import static org.thingsboard.server.common.data.DataConstants.DEVICE; |
|
|
|
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE; |
|
|
|
import static org.thingsboard.server.msa.prototypes.DevicePrototypes.defaultGatewayPrototype; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
Device gatewayDevice; |
|
|
|
MqttClient mqttClient; |
|
|
|
Device createdDevice; |
|
|
|
MqttMessageListener listener; |
|
|
|
private Device gatewayDevice; |
|
|
|
private MqttClient mqttClient; |
|
|
|
private Device createdDevice; |
|
|
|
private MqttMessageListener listener; |
|
|
|
private JsonParser jsonParser = new JsonParser(); |
|
|
|
|
|
|
|
@Before |
|
|
|
@BeforeMethod |
|
|
|
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()); |
|
|
|
testRestClient.login("tenant@thingsboard.org", "tenant"); |
|
|
|
gatewayDevice = testRestClient.postDevice("", defaultGatewayPrototype()); |
|
|
|
DeviceCredentials gatewayDeviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(gatewayDevice.getId()); |
|
|
|
|
|
|
|
this.listener = new MqttMessageListener(); |
|
|
|
this.mqttClient = getMqttClient(gatewayDeviceCredentials.get(), listener); |
|
|
|
this.mqttClient = getMqttClient(gatewayDeviceCredentials, 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()); |
|
|
|
@AfterMethod |
|
|
|
public void removeGateway() { |
|
|
|
testRestClient.deleteDeviceIfExists(this.gatewayDevice.getId()); |
|
|
|
testRestClient.deleteDeviceIfExists(this.createdDevice.getId()); |
|
|
|
this.listener = null; |
|
|
|
this.mqttClient = null; |
|
|
|
this.createdDevice = null; |
|
|
|
@ -95,40 +102,38 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
log.info("Received telemetry: {}", actualLatestTelemetry); |
|
|
|
wsClient.closeBlocking(); |
|
|
|
|
|
|
|
Assert.assertEquals(4, actualLatestTelemetry.getData().size()); |
|
|
|
Assert.assertEquals(Sets.newHashSet("booleanKey", "stringKey", "doubleKey", "longKey"), |
|
|
|
actualLatestTelemetry.getLatestValues().keySet()); |
|
|
|
assertThat(actualLatestTelemetry.getData()).hasSize(4); |
|
|
|
assertThat(actualLatestTelemetry.getLatestValues().keySet()).containsOnlyOnceElementsOf(Arrays.asList("booleanKey", "stringKey", "doubleKey", "longKey")); |
|
|
|
|
|
|
|
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))); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("booleanKey").get(1)).isEqualTo(Boolean.TRUE.toString()); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("stringKey").get(1)).isEqualTo("value1"); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("doubleKey").get(1)).isEqualTo(Double.toString(42.0)); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("longKey").get(1)).isEqualTo(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()); |
|
|
|
assertThat(actualLatestTelemetry.getData()).hasSize(4); |
|
|
|
assertThat(actualLatestTelemetry.getLatestValues().keySet()).containsOnlyOnceElementsOf(Arrays.asList("booleanKey", "stringKey", "doubleKey", "longKey")); |
|
|
|
|
|
|
|
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))); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("booleanKey").get(1)).isEqualTo(Boolean.TRUE.toString()); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("stringKey").get(1)).isEqualTo("value1"); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("doubleKey").get(1)).isEqualTo(Double.toString(42.0)); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("longKey").get(1)).isEqualTo(Long.toString(73)); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void publishAttributeUpdateToServer() throws Exception { |
|
|
|
Optional<DeviceCredentials> createdDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); |
|
|
|
Assert.assertTrue(createdDeviceCredentials.isPresent()); |
|
|
|
testRestClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); |
|
|
|
|
|
|
|
WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); |
|
|
|
JsonObject clientAttributes = new JsonObject(); |
|
|
|
clientAttributes.addProperty("attr1", "value1"); |
|
|
|
@ -142,20 +147,18 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
log.info("Received attributes: {}", actualLatestTelemetry); |
|
|
|
wsClient.closeBlocking(); |
|
|
|
|
|
|
|
Assert.assertEquals(4, actualLatestTelemetry.getData().size()); |
|
|
|
Assert.assertEquals(Sets.newHashSet("attr1", "attr2", "attr3", "attr4"), |
|
|
|
actualLatestTelemetry.getLatestValues().keySet()); |
|
|
|
assertThat(actualLatestTelemetry.getData()).hasSize(4); |
|
|
|
assertThat(actualLatestTelemetry.getLatestValues().keySet()).containsOnlyOnceElementsOf(Arrays.asList("attr1", "attr2", "attr3", "attr4")); |
|
|
|
|
|
|
|
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))); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("attr1").get(1)).isEqualTo("value1"); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("attr2").get(1)).isEqualTo(Boolean.TRUE.toString()); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("attr3").get(1)).isEqualTo(Double.toString(42.0)); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("attr4").get(1)).isEqualTo(Long.toString(73)); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void responseDataOnAttributesRequestCheck() throws Exception { |
|
|
|
Optional<DeviceCredentials> createdDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); |
|
|
|
Assert.assertTrue(createdDeviceCredentials.isPresent()); |
|
|
|
testRestClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); |
|
|
|
JsonObject sharedAttributes = new JsonObject(); |
|
|
|
sharedAttributes.addProperty("attr1", "value1"); |
|
|
|
sharedAttributes.addProperty("attr2", true); |
|
|
|
@ -163,11 +166,8 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
sharedAttributes.addProperty("attr4", 73); |
|
|
|
|
|
|
|
mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); |
|
|
|
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()); |
|
|
|
|
|
|
|
testRestClient.postTelemetryAttribute(DataConstants.DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); |
|
|
|
var event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
JsonObject requestData = new JsonObject(); |
|
|
|
@ -181,8 +181,8 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
JsonObject responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); |
|
|
|
Assert.assertTrue(responseData.has("value")); |
|
|
|
Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("value").getAsString()); |
|
|
|
assertThat(responseData.has("value")).isTrue(); |
|
|
|
assertThat(responseData.get("value").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString()); |
|
|
|
|
|
|
|
requestData = new JsonObject(); |
|
|
|
requestData.addProperty("id", 1); |
|
|
|
@ -198,9 +198,9 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); |
|
|
|
|
|
|
|
Assert.assertTrue(responseData.has("values")); |
|
|
|
Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("values").getAsJsonObject().get("attr1").getAsString()); |
|
|
|
Assert.assertEquals(sharedAttributes.get("attr2").getAsString(), responseData.get("values").getAsJsonObject().get("attr2").getAsString()); |
|
|
|
assertThat(responseData.has("values")).isTrue(); |
|
|
|
assertThat(responseData.get("values").getAsJsonObject().get("attr1").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString()); |
|
|
|
assertThat(responseData.get("values").getAsJsonObject().get("attr2").getAsString()).isEqualTo(sharedAttributes.get("attr2").getAsString()); |
|
|
|
|
|
|
|
requestData = new JsonObject(); |
|
|
|
requestData.addProperty("id", 1); |
|
|
|
@ -216,9 +216,9 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); |
|
|
|
|
|
|
|
Assert.assertTrue(responseData.has("values")); |
|
|
|
Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("values").getAsJsonObject().get("attr1").getAsString()); |
|
|
|
Assert.assertEquals(1, responseData.get("values").getAsJsonObject().entrySet().size()); |
|
|
|
assertThat(responseData.has("values")).isTrue(); |
|
|
|
assertThat(responseData.get("values").getAsJsonObject().get("attr1").getAsString()).isEqualTo(sharedAttributes.get("attr1").getAsString()); |
|
|
|
assertThat(responseData.get("values").getAsJsonObject().entrySet()).hasSize(1); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -237,11 +237,9 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
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)); |
|
|
|
assertThat(actualLatestTelemetry.getData()).hasSize(1); |
|
|
|
assertThat(actualLatestTelemetry.getLatestValues().keySet()).containsOnly("clientAttr"); |
|
|
|
assertThat(actualLatestTelemetry.getDataValuesByKey("clientAttr").get(1)).isEqualTo(clientAttributeValue); |
|
|
|
|
|
|
|
// Add a new shared attribute
|
|
|
|
JsonObject sharedAttributes = new JsonObject(); |
|
|
|
@ -251,16 +249,12 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
// Subscribe for attribute update event
|
|
|
|
mqttClient.on("v1/gateway/attributes", listener, MqttQoS.AT_LEAST_ONCE).get(); |
|
|
|
|
|
|
|
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()); |
|
|
|
testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); |
|
|
|
MqttEvent sharedAttributeEvent = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
// Catch attribute update event
|
|
|
|
Assert.assertNotNull(sharedAttributeEvent); |
|
|
|
Assert.assertEquals("v1/gateway/attributes", sharedAttributeEvent.getTopic()); |
|
|
|
assertThat(sharedAttributeEvent).isNotNull(); |
|
|
|
assertThat(sharedAttributeEvent.getTopic()).isEqualTo("v1/gateway/attributes"); |
|
|
|
|
|
|
|
// Subscribe to attributes response
|
|
|
|
mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); |
|
|
|
@ -288,15 +282,11 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
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()); |
|
|
|
testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); |
|
|
|
|
|
|
|
MqttEvent event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
Assert.assertEquals(sharedAttributeValue, |
|
|
|
mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()); |
|
|
|
assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()) |
|
|
|
.isEqualTo(sharedAttributeValue); |
|
|
|
|
|
|
|
// Update the shared attribute value
|
|
|
|
JsonObject updatedSharedAttributes = new JsonObject(); |
|
|
|
@ -307,15 +297,10 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
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()); |
|
|
|
|
|
|
|
testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(updatedSharedAttributes.toString())); |
|
|
|
event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
Assert.assertEquals(updatedSharedAttributeValue, |
|
|
|
mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()); |
|
|
|
assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()) |
|
|
|
.isEqualTo(updatedSharedAttributeValue); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -331,14 +316,11 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
serverRpcPayload.addProperty("method", "getValue"); |
|
|
|
serverRpcPayload.addProperty("params", true); |
|
|
|
ListeningExecutorService service = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName(getClass().getSimpleName()))); |
|
|
|
ListenableFuture<ResponseEntity> future = service.submit(() -> { |
|
|
|
ListenableFuture<JsonNode> future = service.submit(() -> { |
|
|
|
try { |
|
|
|
return restClient.getRestTemplate() |
|
|
|
.postForEntity(HTTPS_URL + "/api/rpc/twoway/{deviceId}", |
|
|
|
mapper.readTree(serverRpcPayload.toString()), String.class, |
|
|
|
createdDevice.getId()); |
|
|
|
return testRestClient.postServerSideRpc(createdDevice.getId(), mapper.readTree(serverRpcPayload.toString())); |
|
|
|
} catch (IOException e) { |
|
|
|
return ResponseEntity.badRequest().build(); |
|
|
|
return null; |
|
|
|
} |
|
|
|
}); |
|
|
|
|
|
|
|
@ -346,19 +328,13 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
MqttEvent requestFromServer = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
service.shutdownNow(); |
|
|
|
|
|
|
|
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(); |
|
|
|
assertThat(requestFromServer).isNotNull(); |
|
|
|
assertThat(requestFromServer.getMessage()).isNotNull(); |
|
|
|
JsonNode requestFromServerJson = JacksonUtil.toJsonNode(requestFromServer.getMessage()); |
|
|
|
assertThat(requestFromServerJson.get("device").asText()).isEqualTo(createdDevice.getName()); |
|
|
|
assertThat(requestFromServerJson.get("data").get("method").asText()).isEqualTo("getValue"); |
|
|
|
assertThat(requestFromServerJson.get("data").get("params").asText()).isEqualTo("true"); |
|
|
|
int requestId = requestFromServerJson.get("data").get("id").asInt(); |
|
|
|
|
|
|
|
JsonObject clientResponse = new JsonObject(); |
|
|
|
clientResponse.addProperty("response", "someResponse"); |
|
|
|
@ -369,16 +345,15 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
// Send a response to the server's RPC request
|
|
|
|
|
|
|
|
mqttClient.publish(gatewayRpcTopic, Unpooled.wrappedBuffer(gatewayResponse.toString().getBytes())).get(); |
|
|
|
ResponseEntity serverResponse = future.get(5 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
Assert.assertTrue(serverResponse.getStatusCode().is2xxSuccessful()); |
|
|
|
Assert.assertEquals(clientResponse.toString(), serverResponse.getBody()); |
|
|
|
JsonNode serverResponse = future.get(5 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
assertThat(serverResponse).isEqualTo(mapper.readTree(clientResponse.toString())); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void deviceCreationAfterDeleted() throws Exception { |
|
|
|
restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + this.createdDevice.getId()); |
|
|
|
Optional<Device> deletedDevice = restClient.getDeviceById(this.createdDevice.getId()); |
|
|
|
Assert.assertTrue(deletedDevice.isEmpty()); |
|
|
|
testRestClient.deleteDevice(this.createdDevice.getId()); |
|
|
|
testRestClient.getDeviceById(this.createdDevice.getId(), HttpStatus.NOT_FOUND.value()); |
|
|
|
this.createdDevice = createDeviceThroughGateway(mqttClient, gatewayDevice); |
|
|
|
} |
|
|
|
|
|
|
|
@ -397,13 +372,13 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
log.info(gatewayAttributesRequest.toString()); |
|
|
|
mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(gatewayAttributesRequest.toString().getBytes())).get(); |
|
|
|
MqttEvent clientAttributeEvent = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); |
|
|
|
Assert.assertNotNull(clientAttributeEvent); |
|
|
|
assertThat(clientAttributeEvent).isNotNull(); |
|
|
|
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()); |
|
|
|
assertThat(responseMessage.get("id").getAsInt()).isEqualTo(messageId); |
|
|
|
assertThat(responseMessage.get("device").getAsString()).isEqualTo(createdDevice.getName()); |
|
|
|
assertThat(responseMessage.entrySet()).hasSize(3); |
|
|
|
assertThat(responseMessage.get("value").getAsString()).isEqualTo(expectedValue); |
|
|
|
} |
|
|
|
|
|
|
|
private Device createDeviceThroughGateway(MqttClient mqttClient, Device gatewayDevice) throws Exception { |
|
|
|
@ -411,24 +386,19 @@ public class MqttGatewayClientTest extends AbstractContainerTest { |
|
|
|
TimeUnit.SECONDS.sleep(30); |
|
|
|
} |
|
|
|
|
|
|
|
String deviceName = "mqtt_device"; |
|
|
|
String deviceName = "mqtt_device" + RandomStringUtils.randomAlphabetic(5); |
|
|
|
mqttClient.publish("v1/gateway/connect", Unpooled.wrappedBuffer(createGatewayConnectPayload(deviceName).toString().getBytes()), MqttQoS.AT_LEAST_ONCE).get(); |
|
|
|
|
|
|
|
if (timeoutMultiplier > 1) { |
|
|
|
TimeUnit.SECONDS.sleep(30); |
|
|
|
} |
|
|
|
|
|
|
|
List<EntityRelation> relations = restClient.findByFrom(gatewayDevice.getId(), RelationTypeGroup.COMMON); |
|
|
|
|
|
|
|
Assert.assertEquals(1, relations.size()); |
|
|
|
List<EntityRelation> relations = testRestClient.findRelationByFrom(gatewayDevice.getId(), RelationTypeGroup.COMMON); |
|
|
|
assertThat(relations).hasSize(1); |
|
|
|
|
|
|
|
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(); |
|
|
|
return testRestClient.getDeviceById(createdDeviceId); |
|
|
|
} |
|
|
|
|
|
|
|
private MqttClient getMqttClient(DeviceCredentials deviceCredentials, MqttMessageListener listener) throws InterruptedException, ExecutionException { |
|
|
|
|