From a3e31c109c30a725cc2fab600f6e84dbd5dacf63 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Thu, 9 Feb 2023 18:49:05 +0200 Subject: [PATCH 1/2] sparkplug: add attributes by Metric`s name, STATE ONLINE --- ...ltDeviceProfileTransportConfiguration.java | 1 - ...ttDeviceProfileTransportConfiguration.java | 3 ++ .../transport/mqtt/MqttTransportHandler.java | 4 -- .../AbstractGatewaySessionHandler.java | 8 ++++ .../session/SparkplugNodeSessionHandler.java | 41 +++++++++++-------- .../sparkplug/SparkplugMessageTypeSate.java | 30 ++++++++++++++ 6 files changed, 64 insertions(+), 23 deletions(-) create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMessageTypeSate.java diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DefaultDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DefaultDeviceProfileTransportConfiguration.java index 0d2251befc..a8765f7fd0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DefaultDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DefaultDeviceProfileTransportConfiguration.java @@ -16,7 +16,6 @@ package org.thingsboard.server.common.data.device.profile; import lombok.Data; -import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceTransportType; @Data diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java index a2cc91ff1d..ead7e1f8e3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/MqttDeviceProfileTransportConfiguration.java @@ -19,6 +19,8 @@ import lombok.Data; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.validation.NoXss; +import java.util.Set; + @Data public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTransportConfiguration { @@ -28,6 +30,7 @@ public class MqttDeviceProfileTransportConfiguration implements DeviceProfileTra private String deviceAttributesTopic = MqttTopics.DEVICE_ATTRIBUTES_TOPIC; private TransportPayloadTypeConfiguration transportPayloadTypeConfiguration; private boolean sparkPlug; + private Set sparkPlugAttributesMetricNames; private boolean sendAckOnValidationException; @Override diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 2953183027..eb76e89bad 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -412,10 +412,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement // TODO break; case DBIRTH: - sparkplugSessionHandler.onTelemetryProto(msgId, sparkplugBProtoDevice, deviceName, sparkplugTopic); - - System.out.println(); - break; case DCMD: case DDATA: sparkplugSessionHandler.onTelemetryProto(msgId, sparkplugBProtoDevice, deviceName, sparkplugTopic); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index 66c5823f11..c9583c12c5 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -180,6 +180,14 @@ public abstract class AbstractGatewaySessionHandler { void deregisterSession(String deviceName) { MqttDeviceAwareSessionContext deviceSessionCtx = devices.remove(deviceName); if (deviceSessionCtx != null) { + if (deviceSessionCtx.isSparkplug()){ +// onDeviceTelemetryProto(int msgId, ByteBuf payload) + // add Metric name: STATE type: String value: OFFLINE +// SparkplugBProto.Payload.Metric metricState = createMetric(SparkplugMessageTypeSate.OFFLINE.name(), +// sparkplugBProtoDevice.getTimestamp(), STATE.name(), MetricDataType.Text); +// sparkplugBProtoDevice.getMetricsList().add(metricState); +// onTelemetryProto(-1, sparkplugBProtoDevice, deviceName, sparkplugTopic); + } deregisterSession(deviceName, deviceSessionCtx); } else { log.debug("[{}] Device [{}] was already removed from the gateway session", sessionId, deviceName); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index ef10c3234e..9af5e71560 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java @@ -51,6 +51,8 @@ import java.util.stream.Collectors; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.DBIRTH; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NBIRTH; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.ONLINE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.fromSparkplugBMetricToKeyValueProto; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.validatedValueByTypeMetric; @@ -102,6 +104,12 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { List msgs = convertToPostTelemetry(sparkplugBProto, topic.getType().name()); if (topic.isType(NBIRTH) || topic.isType(DBIRTH)) { try { + // add Msg Telemetry: key STATE type: String value: ONLINE ts: sparkplugBProto.getTimestamp() + TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); + keyValueProtoBuilder.setKey(STATE.name()); + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); + keyValueProtoBuilder.setStringV(ONLINE.name()); + msgs.add(postTelemetryMsgCreated(keyValueProtoBuilder.build(), sparkplugBProto.getTimestamp())); contextListenableFuture.get().setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); } catch (InterruptedException | ExecutionException e) { log.error("Failed add Metrics. MessageType *BIRTH.", e); @@ -182,30 +190,16 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { topicTypeName + " " + protoMetric.getName() : protoMetric.getName(); Optional keyValueProtoOpt = fromSparkplugBMetricToKeyValueProto(key, protoMetric); if (keyValueProtoOpt.isPresent()) { - List result = new ArrayList<>(); - result.add(keyValueProtoOpt.get()); - TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); - TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); - builder.setTs(ts); - builder.addAllKv(result); - request.addTsKvList(builder.build()); - msgs.add(request.build()); + msgs.add(postTelemetryMsgCreated(keyValueProtoOpt.get(), ts)); } } if (DBIRTH.name().equals(topicTypeName)) { - List result = new ArrayList<>(); TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); keyValueProtoBuilder.setKey(topicTypeName + " " + "seq"); keyValueProtoBuilder.setType(TransportProtos.KeyValueType.LONG_V); keyValueProtoBuilder.setLongV(sparkplugBProto.getSeq()); - result.add(keyValueProtoBuilder.build()); - TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); - TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); - builder.setTs(sparkplugBProto.getTimestamp()); - builder.addAllKv(result); - request.addTsKvList(builder.build()); - msgs.add(request.build()); + msgs.add(postTelemetryMsgCreated(keyValueProtoBuilder.build(), sparkplugBProto.getTimestamp())); } return msgs; } catch (IllegalStateException | JsonSyntaxException | ThingsboardException e) { @@ -233,11 +227,11 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { return Optional.of(getPayloadAdaptor().createMqttPublishMsg(deviceSessionCtx, sparkplugTopic, payloadInBytes)); } else { log.trace("DeviceId: [{}] tenantId: [{}] sessionId:[{}] Failed to convert device attributes [{}] response to MQTT sparkplug msg", - deviceSessionCtx.getDeviceInfo().getDeviceId(), deviceSessionCtx.getDeviceInfo().getTenantId(), sessionId, tsKvProto.getKv()); + deviceSessionCtx.getDeviceInfo().getDeviceId(), deviceSessionCtx.getDeviceInfo().getTenantId(), sessionId, tsKvProto.getKv()); } } catch (Exception e) { log.trace("DeviceId: [{}] tenantId: [{}] sessionId:[{}] Failed to convert device attributes response to MQTT sparkplug msg", - deviceSessionCtx.getDeviceInfo().getDeviceId(), deviceSessionCtx.getDeviceInfo().getTenantId(), sessionId, e); + deviceSessionCtx.getDeviceInfo().getDeviceId(), deviceSessionCtx.getDeviceInfo().getTenantId(), sessionId, e); return Optional.empty(); } return Optional.empty(); @@ -248,4 +242,15 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { return new SparkplugDeviceSessionContext(this, msg.getDeviceInfo(), msg.getDeviceProfile(), mqttQoSMap, transportService); } + protected TransportProtos.PostTelemetryMsg postTelemetryMsgCreated (TransportProtos.KeyValueProto keyValueProto, long ts) { + List result = new ArrayList<>(); + result.add(keyValueProto); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(ts); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + return request.build(); + } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMessageTypeSate.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMessageTypeSate.java new file mode 100644 index 0000000000..f9a18adbd3 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/sparkplug/SparkplugMessageTypeSate.java @@ -0,0 +1,30 @@ +/** + * Copyright © 2016-2022 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.transport.mqtt.util.sparkplug; + +public enum SparkplugMessageTypeSate { + /** + * The EoN node should examine the payload of this + * message to ensure that it is a value of “ONLINE” + */ + OFFLINE, + /** + * If the value is “OFFLINE”, this indicates the Primary Application + * has lost its MQTT Session to this particular MQTT Server. + */ + ONLINE + +} From 6545af2282ec29b05d4030678738b26621d267ab Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Fri, 10 Feb 2023 15:26:27 +0200 Subject: [PATCH 2/2] sparkplug: STATE ONLINE/OFFLINE --- .../transport/mqtt/MqttTransportHandler.java | 11 ++++- .../AbstractGatewaySessionHandler.java | 43 ++++++++++++++----- .../session/SparkplugNodeSessionHandler.java | 19 +------- 3 files changed, 45 insertions(+), 28 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index eb76e89bad..5a8927e184 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -91,6 +91,7 @@ import java.security.cert.Certificate; import java.security.cert.X509Certificate; import java.util.ArrayList; import java.util.Collections; +import java.util.Date; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -110,6 +111,7 @@ import static io.netty.handler.codec.mqtt.MqttQoS.AT_MOST_ONCE; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_CLOSED; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NDEATH; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.OFFLINE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicUtil.parseTopicPublish; /** @@ -1124,7 +1126,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(deviceSessionCtx.getSessionInfo(), SESSION_EVENT_MSG_CLOSED, null); transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); if (gatewaySessionHandler != null) { - gatewaySessionHandler.onGatewayDisconnect(); + gatewaySessionHandler.onDevicesDisconnect(); + } + if (sparkplugSessionHandler != null) { + // add Msg Telemetry node: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() + sparkplugSessionHandler.stateSparkplugtSendOnTelemetry(deviceSessionCtx.getSessionInfo(), + deviceSessionCtx.getDeviceInfo().getDeviceName(), OFFLINE, new Date().getTime()); + sparkplugSessionHandler.onDevicesDisconnect(); } deviceSessionCtx.setDisconnected(); } @@ -1224,6 +1232,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); + transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); deviceSessionCtx.getChannel().close(); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index c9583c12c5..ba8036bb9d 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java @@ -54,9 +54,12 @@ import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor; import org.thingsboard.server.transport.mqtt.util.ReturnCode; +import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate; import javax.annotation.Nullable; +import java.util.ArrayList; import java.util.Collections; +import java.util.Date; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -72,6 +75,8 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_OPEN; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG; import static org.thingsboard.server.common.transport.service.DefaultTransportService.SUBSCRIBE_TO_RPC_ASYNC_MSG; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; +import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.OFFLINE; /** * Created by ashvayka on 19.01.17. @@ -157,7 +162,7 @@ public abstract class AbstractGatewaySessionHandler { } } - public void onGatewayDisconnect() { + public void onDevicesDisconnect() { devices.forEach(this::deregisterSession); } @@ -180,14 +185,6 @@ public abstract class AbstractGatewaySessionHandler { void deregisterSession(String deviceName) { MqttDeviceAwareSessionContext deviceSessionCtx = devices.remove(deviceName); if (deviceSessionCtx != null) { - if (deviceSessionCtx.isSparkplug()){ -// onDeviceTelemetryProto(int msgId, ByteBuf payload) - // add Metric name: STATE type: String value: OFFLINE -// SparkplugBProto.Payload.Metric metricState = createMetric(SparkplugMessageTypeSate.OFFLINE.name(), -// sparkplugBProtoDevice.getTimestamp(), STATE.name(), MetricDataType.Text); -// sparkplugBProtoDevice.getMetricsList().add(metricState); -// onTelemetryProto(-1, sparkplugBProtoDevice, deviceName, sparkplugTopic); - } deregisterSession(deviceName, deviceSessionCtx); } else { log.debug("[{}] Device [{}] was already removed from the gateway session", sessionId, deviceName); @@ -403,10 +400,21 @@ public abstract class AbstractGatewaySessionHandler { } } - protected void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { + public void processPostTelemetryMsg(MqttDeviceAwareSessionContext deviceCtx, TransportProtos.PostTelemetryMsg postTelemetryMsg, String deviceName, int msgId) { transportService.process(deviceCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(channel, deviceName, msgId, postTelemetryMsg)); } + public TransportProtos.PostTelemetryMsg postTelemetryMsgCreated(TransportProtos.KeyValueProto keyValueProto, long ts) { + List result = new ArrayList<>(); + result.add(keyValueProto); + TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); + TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); + builder.setTs(ts); + builder.addAllKv(result); + request.addTsKvList(builder.build()); + return request.build(); + } + private void onDeviceClaimJson(int msgId, ByteBuf payload) throws AdaptorException { JsonElement json = JsonMqttAdaptor.validateJsonPayload(sessionId, payload); if (json.isJsonObject()) { @@ -719,12 +727,26 @@ public abstract class AbstractGatewaySessionHandler { } private void deregisterSession(String deviceName, MqttDeviceAwareSessionContext deviceSessionCtx) { + if (this.deviceSessionCtx.isSparkplug()){ + // add Msg Telemetry: key STATE type: String value: OFFLINE ts: sparkplugBProto.getTimestamp() + stateSparkplugtSendOnTelemetry (deviceSessionCtx.getSessionInfo(), + deviceSessionCtx.getDeviceInfo().getDeviceName(), OFFLINE, new Date().getTime()); + } transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); transportService.process(deviceSessionCtx.getSessionInfo(), SESSION_EVENT_MSG_CLOSED, null); System.out.println("Removed device " + deviceName + " from the gateway session"); log.debug("[{}] Removed device [{}] from the gateway session", sessionId, deviceName); } + public void stateSparkplugtSendOnTelemetry (TransportProtos.SessionInfoProto sessionInfo, String deviceName, SparkplugMessageTypeSate typeSate, long ts) { + TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); + keyValueProtoBuilder.setKey(STATE.name()); + keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); + keyValueProtoBuilder.setStringV(typeSate.name()); + TransportProtos.PostTelemetryMsg postTelemetryMsg = postTelemetryMsgCreated(keyValueProtoBuilder.build(), ts); + transportService.process(sessionInfo, postTelemetryMsg, getPubAckCallback(channel, deviceName, -1, postTelemetryMsg)); + } + private TransportServiceCallback getPubAckCallback(final ChannelHandlerContext ctx, final String deviceName, final int msgId, final T msg) { return new TransportServiceCallback() { @Override @@ -742,4 +764,5 @@ public abstract class AbstractGatewaySessionHandler { } }; } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index 9af5e71560..c52c5e33a8 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java @@ -51,7 +51,6 @@ import java.util.stream.Collectors; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.DBIRTH; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.NBIRTH; -import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageTypeSate.ONLINE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.createMetric; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetricUtil.fromSparkplugBMetricToKeyValueProto; @@ -105,11 +104,8 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { if (topic.isType(NBIRTH) || topic.isType(DBIRTH)) { try { // add Msg Telemetry: key STATE type: String value: ONLINE ts: sparkplugBProto.getTimestamp() - TransportProtos.KeyValueProto.Builder keyValueProtoBuilder = TransportProtos.KeyValueProto.newBuilder(); - keyValueProtoBuilder.setKey(STATE.name()); - keyValueProtoBuilder.setType(TransportProtos.KeyValueType.STRING_V); - keyValueProtoBuilder.setStringV(ONLINE.name()); - msgs.add(postTelemetryMsgCreated(keyValueProtoBuilder.build(), sparkplugBProto.getTimestamp())); + stateSparkplugtSendOnTelemetry(contextListenableFuture.get().getSessionInfo(), deviceName, ONLINE, + sparkplugBProto.getTimestamp()); contextListenableFuture.get().setDeviceBirthMetrics(sparkplugBProto.getMetricsList()); } catch (InterruptedException | ExecutionException e) { log.error("Failed add Metrics. MessageType *BIRTH.", e); @@ -242,15 +238,4 @@ public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { return new SparkplugDeviceSessionContext(this, msg.getDeviceInfo(), msg.getDeviceProfile(), mqttQoSMap, transportService); } - protected TransportProtos.PostTelemetryMsg postTelemetryMsgCreated (TransportProtos.KeyValueProto keyValueProto, long ts) { - List result = new ArrayList<>(); - result.add(keyValueProto); - TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); - TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder(); - builder.setTs(ts); - builder.addAllKv(result); - request.addTsKvList(builder.build()); - return request.build(); - } - }