|
|
|
@ -16,9 +16,7 @@ |
|
|
|
package org.thingsboard.server.transport.mqtt.sparkplug.connection; |
|
|
|
|
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import com.google.protobuf.ByteString; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.eclipse.paho.client.mqttv3.MqttClient; |
|
|
|
import org.eclipse.paho.mqttv5.client.IMqttToken; |
|
|
|
import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; |
|
|
|
import org.eclipse.paho.mqttv5.common.MqttMessage; |
|
|
|
@ -26,33 +24,25 @@ import org.eclipse.paho.mqttv5.common.packet.MqttConnAck; |
|
|
|
import org.eclipse.paho.mqttv5.common.packet.MqttReturnCode; |
|
|
|
import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; |
|
|
|
import org.junit.Assert; |
|
|
|
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; |
|
|
|
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|
|
|
import org.thingsboard.server.common.data.Device; |
|
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.LongDataEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
|
import org.thingsboard.server.gen.transport.mqtt.SparkplugBProto; |
|
|
|
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; |
|
|
|
import org.thingsboard.server.transport.mqtt.mqttv5.MqttV5TestClient; |
|
|
|
import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; |
|
|
|
import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; |
|
|
|
import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; |
|
|
|
import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil; |
|
|
|
|
|
|
|
import java.io.ByteArrayOutputStream; |
|
|
|
import java.io.DataOutputStream; |
|
|
|
import java.io.IOException; |
|
|
|
import java.io.ObjectOutputStream; |
|
|
|
import java.nio.ByteBuffer; |
|
|
|
import java.util.Calendar; |
|
|
|
import java.util.Date; |
|
|
|
|
|
|
|
import java.util.Optional; |
|
|
|
import java.util.concurrent.ExecutorService; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.HashSet; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.atomic.AtomicReference; |
|
|
|
|
|
|
|
import static org.awaitility.Awaitility.await; |
|
|
|
import static org.eclipse.paho.mqttv5.common.packet.MqttWireMessage.MESSAGE_TYPE_CONNACK; |
|
|
|
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int32; |
|
|
|
import static org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType.Int64; |
|
|
|
|
|
|
|
/** |
|
|
|
@ -67,17 +57,15 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr |
|
|
|
|
|
|
|
protected void processClientWithCorrectNodeAccessTokenWithNdeathTest() throws Exception { |
|
|
|
long ts = calendar.getTimeInMillis()-PUBLISH_TS_DELTA_MS; |
|
|
|
int value = bdSeq; |
|
|
|
long value = bdSeq = 0; |
|
|
|
MetricDataType metricDataType = Int64; |
|
|
|
TsKvEntry expectedTsKvEntryOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, Integer.toUnsignedLong(value))); |
|
|
|
TsKvEntry tsKvEntryBdSecOriginal = new BasicTsKvEntry(ts, new LongDataEntry(keysBdSeq, value)); |
|
|
|
|
|
|
|
SparkplugBProto.Payload.Builder deathPayload = SparkplugBProto.Payload.newBuilder() |
|
|
|
.setTimestamp(calendar.getTimeInMillis()); |
|
|
|
deathPayload.addMetrics(createMetric(expectedTsKvEntryOriginal, metricDataType)); |
|
|
|
|
|
|
|
byte[] deathBytes = deathPayload.build().toByteArray(); |
|
|
|
deathPayload.addMetrics(createMetric(tsKvEntryBdSecOriginal, metricDataType)); |
|
|
|
|
|
|
|
MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNdeath(deathBytes); |
|
|
|
MqttWireMessage response = clientWithCorrectNodeAccessTokenWithNDEATH(deathPayload.build().toByteArray()); |
|
|
|
|
|
|
|
Assert.assertEquals(MESSAGE_TYPE_CONNACK, response.getType()); |
|
|
|
|
|
|
|
@ -86,10 +74,9 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr |
|
|
|
Assert.assertEquals(MqttReturnCode.RETURN_CODE_SUCCESS, connAckMsg.getReturnCode()); |
|
|
|
|
|
|
|
String keys = SparkplugMessageType.NDEATH.name() + " " + keysBdSeq; |
|
|
|
TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(value))); |
|
|
|
ListenableFuture<Optional<TsKvEntry>> future = tsService.findLatest(tenantId, savedGateway.getId(), keys); |
|
|
|
AtomicReference<ListenableFuture<Optional<TsKvEntry>>> finalFuture = new AtomicReference<>(future); |
|
|
|
await("Failed Post Telemetry node proto payload. SparkplugMessageType NDEATH") |
|
|
|
TsKvEntry expectedTsKvEntry = new BasicTsKvEntry(ts, new LongDataEntry(keys, value)); |
|
|
|
AtomicReference<ListenableFuture<Optional<TsKvEntry>>> finalFuture = new AtomicReference<>(); |
|
|
|
await(alias + SparkplugMessageType.NDEATH.name()) |
|
|
|
.atMost(40, TimeUnit.SECONDS) |
|
|
|
.until(() -> { |
|
|
|
finalFuture.set(tsService.findLatest(tenantId, savedGateway.getId(), keys)); |
|
|
|
@ -97,17 +84,58 @@ public abstract class AbstractMqttV5ClientSparkplugConnectionTest extends Abstr |
|
|
|
}); |
|
|
|
TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); |
|
|
|
Assert.assertEquals(expectedTsKvEntry, actualTsKvEntry); |
|
|
|
client.disconnect(); |
|
|
|
} |
|
|
|
|
|
|
|
protected void processClientWithCorrectAccessTokenCreatedDevices(int cntDevices) throws Exception { |
|
|
|
processClientWithCorrectNodeAccess(); |
|
|
|
long ts = calendar.getTimeInMillis(); |
|
|
|
MetricDataType metricDataType = Int32; |
|
|
|
Set<String> deviceIds = new HashSet<>(); |
|
|
|
String keys = "Device Metric int32"; |
|
|
|
int valueDeviceInt32 = 1024; |
|
|
|
TsKvEntry expectedTsKvEntryDeviceInt32 = new BasicTsKvEntry(ts, new LongDataEntry(keys, Integer.toUnsignedLong(valueDeviceInt32))); |
|
|
|
SparkplugBProto.Payload.Metric metric = createMetric(expectedTsKvEntryDeviceInt32, metricDataType); |
|
|
|
for (int i=0; i < cntDevices; i++ ) { |
|
|
|
SparkplugBProto.Payload.Builder payloadBirthDevice = SparkplugBProto.Payload.newBuilder() |
|
|
|
.setTimestamp(calendar.getTimeInMillis()) |
|
|
|
.setSeq(getSeqNum()); |
|
|
|
String deviceName = deviceId + "_" + i; |
|
|
|
|
|
|
|
payloadBirthDevice.addMetrics(metric); |
|
|
|
if (client.isConnected()) { |
|
|
|
client.publish(NAMESPACE + "/" + groupId + "/" + SparkplugMessageType.DBIRTH.name() + "/" + edgeNode + "/" + deviceName, |
|
|
|
payloadBirthDevice.build().toByteArray(), 0, false); |
|
|
|
deviceIds.add(deviceName); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
Assert.assertEquals(cntDevices, deviceIds.size()); |
|
|
|
|
|
|
|
for (String deviceName: deviceIds) { |
|
|
|
AtomicReference<Device> device = new AtomicReference<>(); |
|
|
|
await(alias + SparkplugMessageType.DBIRTH.name()) |
|
|
|
.atMost(40, TimeUnit.SECONDS) |
|
|
|
.until(() -> { |
|
|
|
device.set(doGet("/api/tenant/devices?deviceName=" + deviceName, Device.class)); |
|
|
|
return device.get() != null; |
|
|
|
}); |
|
|
|
Assert.assertEquals(deviceName, device.get().getName()); |
|
|
|
AtomicReference<ListenableFuture<Optional<TsKvEntry>>> finalFuture = new AtomicReference<>(); |
|
|
|
await(alias + SparkplugMessageType.NDEATH.name()) |
|
|
|
.atMost(40, TimeUnit.SECONDS) |
|
|
|
.until(() -> { |
|
|
|
finalFuture.set(tsService.findLatest(tenantId, device.get().getId(), keys)); |
|
|
|
return finalFuture.get().get().isPresent(); |
|
|
|
}); |
|
|
|
TsKvEntry actualTsKvEntry = finalFuture.get().get().get(); |
|
|
|
Assert.assertEquals(expectedTsKvEntryDeviceInt32, actualTsKvEntry); |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
private MqttWireMessage clientWithCorrectNodeAccessTokenWithNdeath(byte[] deathBytes) throws Exception { |
|
|
|
private MqttWireMessage clientWithCorrectNodeAccessTokenWithNDEATH(byte[] deathBytes) throws Exception { |
|
|
|
this.client = new MqttV5TestClient(); |
|
|
|
MqttConnectionOptions options = new MqttConnectionOptions(); |
|
|
|
options.setUserName(gatewayAccessToken); |
|
|
|
|