Browse Source

Mqtt transport implementation: POST telemetry, attributes

pull/1166/head
Andrew Shvayka 8 years ago
parent
commit
2637babfe3
  1. 8
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 8
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  3. 7
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java
  4. 306
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  5. 36
      application/src/main/java/org/thingsboard/server/actors/device/RuleEngineQueuePutAckMsg.java
  6. 8
      application/src/main/java/org/thingsboard/server/actors/device/SessionInfo.java
  7. 40
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
  8. 14
      application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
  9. 16
      application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java
  10. 17
      application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java
  11. 9
      application/src/main/java/org/thingsboard/server/service/transport/RuleEngineTransportService.java
  12. 8
      application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java
  13. 16
      application/src/main/java/org/thingsboard/server/service/transport/msg/TransportToDeviceActorMsgWrapper.java
  14. 3
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  15. 52
      common/message/src/main/java/org/thingsboard/server/common/msg/core/BasicActorSystemToDeviceSessionActorMsg.java
  16. 36
      common/message/src/main/java/org/thingsboard/server/common/msg/timeout/DeviceActorQueueTimeoutMsg.java
  17. 8
      common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaConsumerTemplate.java
  18. 12
      common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java
  19. 2
      common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java
  20. 8
      common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java
  21. 8
      common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java
  22. 2
      common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgProcessor.java
  23. 1
      common/transport/src/main/java/org/thingsboard/server/common/transport/TransportAdaptor.java
  24. 11
      common/transport/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  25. 180
      common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java
  26. 8
      common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java
  27. 36
      common/transport/src/main/proto/transport.proto
  28. 48
      transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java
  29. 302
      transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  30. 61
      transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java
  31. 13
      transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java
  32. 4
      transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
  33. 33
      transport/mqtt-transport/src/main/java/org/thingsboard/server/mqtt/service/MqttTransportService.java
  34. 8
      transport/mqtt-transport/src/main/java/org/thingsboard/server/mqtt/service/ToRuleEngineMsgEncoder.java

8
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -198,10 +198,6 @@ public class ActorSystemContext {
@Getter @Getter
private MailService mailService; private MailService mailService;
@Autowired
@Getter
private MsgQueueService msgQueueService;
@Autowired @Autowired
@Getter @Getter
private DeviceStateService deviceStateService; private DeviceStateService deviceStateService;
@ -267,10 +263,6 @@ public class ActorSystemContext {
@Setter @Setter
private ActorRef appActor; private ActorRef appActor;
@Getter
@Setter
private ActorRef sessionManagerActor;
@Getter @Getter
@Setter @Setter
private ActorRef statsActor; private ActorRef statsActor;

8
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -38,7 +38,6 @@ import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.aware.TenantAwareMsg; import org.thingsboard.server.common.msg.aware.TenantAwareMsg;
import org.thingsboard.server.common.msg.cluster.SendToClusterMsg; import org.thingsboard.server.common.msg.cluster.SendToClusterMsg;
import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.common.msg.cluster.ServerAddress;
import org.thingsboard.server.common.msg.core.BasicActorSystemToDeviceSessionActorMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg; import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
@ -113,19 +112,12 @@ public class AppActor extends RuleChainManagerActor {
case REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG: case REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG:
onToDeviceActorMsg((TenantAwareMsg) msg); onToDeviceActorMsg((TenantAwareMsg) msg);
break; break;
case ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG:
onToDeviceSessionMsg((BasicActorSystemToDeviceSessionActorMsg) msg);
break;
default: default:
return false; return false;
} }
return true; return true;
} }
private void onToDeviceSessionMsg(BasicActorSystemToDeviceSessionActorMsg msg) {
systemContext.getSessionManagerActor().tell(msg, self());
}
private void onPossibleClusterMsg(SendToClusterMsg msg) { private void onPossibleClusterMsg(SendToClusterMsg msg) {
Optional<ServerAddress> address = systemContext.getRoutingService().resolveById(msg.getEntityId()); Optional<ServerAddress> address = systemContext.getRoutingService().resolveById(msg.getEntityId());
if (address.isPresent()) { if (address.isPresent()) {

7
application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java

@ -27,7 +27,6 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; import org.thingsboard.server.common.msg.cluster.ClusterEventMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorClientSideRpcTimeoutMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorClientSideRpcTimeoutMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorQueueTimeoutMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg;
import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg;
import org.thingsboard.server.service.rpc.ToServerRpcResponseActorMsg; import org.thingsboard.server.service.rpc.ToServerRpcResponseActorMsg;
@ -74,12 +73,6 @@ public class DeviceActor extends ContextAwareActor {
case DEVICE_ACTOR_CLIENT_SIDE_RPC_TIMEOUT_MSG: case DEVICE_ACTOR_CLIENT_SIDE_RPC_TIMEOUT_MSG:
processor.processClientSideRpcTimeout(context(), (DeviceActorClientSideRpcTimeoutMsg) msg); processor.processClientSideRpcTimeout(context(), (DeviceActorClientSideRpcTimeoutMsg) msg);
break; break;
case DEVICE_ACTOR_QUEUE_TIMEOUT_MSG:
processor.processQueueTimeout(context(), (DeviceActorQueueTimeoutMsg) msg);
break;
case RULE_ENGINE_QUEUE_PUT_ACK_MSG:
processor.processQueueAck(context(), (RuleEngineQueuePutAckMsg) msg);
break;
default: default:
return false; return false;
} }

306
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -16,7 +16,6 @@
package org.thingsboard.server.actors.device; package org.thingsboard.server.actors.device;
import akka.actor.ActorContext; import akka.actor.ActorContext;
import akka.actor.ActorRef;
import akka.event.LoggingAdapter; import akka.event.LoggingAdapter;
import com.datastax.driver.core.utils.UUIDs; import com.datastax.driver.core.utils.UUIDs;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
@ -46,29 +45,15 @@ import org.thingsboard.server.common.msg.cluster.ClusterEventMsg;
import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.common.msg.cluster.ServerAddress;
import org.thingsboard.server.common.msg.core.ActorSystemToDeviceSessionActorMsg; import org.thingsboard.server.common.msg.core.ActorSystemToDeviceSessionActorMsg;
import org.thingsboard.server.common.msg.core.AttributesUpdateNotification; import org.thingsboard.server.common.msg.core.AttributesUpdateNotification;
import org.thingsboard.server.common.msg.core.AttributesUpdateRequest;
import org.thingsboard.server.common.msg.core.BasicActorSystemToDeviceSessionActorMsg;
import org.thingsboard.server.common.msg.core.BasicCommandAckResponse;
import org.thingsboard.server.common.msg.core.BasicGetAttributesResponse;
import org.thingsboard.server.common.msg.core.BasicStatusCodeResponse;
import org.thingsboard.server.common.msg.core.GetAttributesRequest;
import org.thingsboard.server.common.msg.core.RuleEngineError; import org.thingsboard.server.common.msg.core.RuleEngineError;
import org.thingsboard.server.common.msg.core.RuleEngineErrorMsg; import org.thingsboard.server.common.msg.core.RuleEngineErrorMsg;
import org.thingsboard.server.common.msg.core.SessionCloseMsg;
import org.thingsboard.server.common.msg.core.SessionCloseNotification;
import org.thingsboard.server.common.msg.core.SessionOpenMsg;
import org.thingsboard.server.common.msg.core.TelemetryUploadRequest;
import org.thingsboard.server.common.msg.core.ToDeviceRpcRequestMsg; import org.thingsboard.server.common.msg.core.ToDeviceRpcRequestMsg;
import org.thingsboard.server.common.msg.core.ToDeviceRpcResponseMsg;
import org.thingsboard.server.common.msg.core.ToServerRpcRequestMsg;
import org.thingsboard.server.common.msg.kv.BasicAttributeKVMsg; import org.thingsboard.server.common.msg.kv.BasicAttributeKVMsg;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
import org.thingsboard.server.common.msg.session.FromDeviceMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.msg.session.SessionType; import org.thingsboard.server.common.msg.session.SessionType;
import org.thingsboard.server.common.msg.session.ToDeviceMsg; import org.thingsboard.server.common.msg.session.ToDeviceMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorClientSideRpcTimeoutMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorClientSideRpcTimeoutMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorQueueTimeoutMsg;
import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
@ -88,9 +73,7 @@ import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.Set; import java.util.Set;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.TimeoutException;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import org.thingsboard.server.gen.transport.TransportProtos.*; import org.thingsboard.server.gen.transport.TransportProtos.*;
@ -192,19 +175,6 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
} }
void processQueueAck(ActorContext context, RuleEngineQueuePutAckMsg msg) {
PendingSessionMsgData data = pendingMsgs.remove(msg.getId());
if (data != null && data.isReplyOnQueueAck()) {
int remainingAcks = data.getAckMsgCount() - 1;
data.setAckMsgCount(remainingAcks);
logger.debug("[{}] Queue put [{}] ack detected. Remaining acks: {}!", deviceId, msg.getId(), remainingAcks);
if (remainingAcks == 0) {
ToDeviceMsg toDeviceMsg = BasicStatusCodeResponse.onSuccess(data.getSessionMsgType(), data.getRequestId());
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(toDeviceMsg, data.getSessionId()), data.getServerAddress());
}
}
}
private void sendPendingRequests(ActorContext context, SessionId sessionId, SessionType type, Optional<ServerAddress> server) { private void sendPendingRequests(ActorContext context, SessionId sessionId, SessionType type, Optional<ServerAddress> server) {
if (!toDeviceRpcPendingMap.isEmpty()) { if (!toDeviceRpcPendingMap.isEmpty()) {
logger.debug("[{}] Pushing {} pending RPC messages to new async session [{}]", deviceId, toDeviceRpcPendingMap.size(), sessionId); logger.debug("[{}] Pushing {} pending RPC messages to new async session [{}]", deviceId, toDeviceRpcPendingMap.size(), sessionId);
@ -239,8 +209,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
body.getMethod(), body.getMethod(),
body.getParams() body.getParams()
); );
ActorSystemToDeviceSessionActorMsg response = new BasicActorSystemToDeviceSessionActorMsg(rpcRequest, sessionId); // ActorSystemToDeviceSessionActorMsg response = new BasicActorSystemToDeviceSessionActorMsg(rpcRequest, sessionId);
sendMsgToSessionActor(response, server); // sendMsgToSessionActor(response, server);
}; };
} }
@ -292,57 +262,25 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
private void handleGetAttributesRequest(ActorContext context, SessionInfoProto sessionInfo, GetAttributeRequestMsg request) { private void handleGetAttributesRequest(ActorContext context, SessionInfoProto sessionInfo, GetAttributeRequestMsg request) {
ListenableFuture<List<AttributeKvEntry>> clientAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.CLIENT_SCOPE, toOptionalSet(request.getClientAttributeNamesList())); ListenableFuture<List<AttributeKvEntry>> clientAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.CLIENT_SCOPE, toOptionalSet(request.getClientAttributeNamesList()));
ListenableFuture<List<AttributeKvEntry>> sharedAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.SHARED_SCOPE, toOptionalSet(request.getSharedAttributeNamesList())); ListenableFuture<List<AttributeKvEntry>> sharedAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.SHARED_SCOPE, toOptionalSet(request.getSharedAttributeNamesList()));
UUID sessionId = new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB());
Futures.addCallback(Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture)), new FutureCallback<List<List<AttributeKvEntry>>>() { int requestId = request.getRequestId();
@Override
public void onSuccess(@Nullable List<List<AttributeKvEntry>> result) {
systemContext.getRuleEngineTransportService().process();
BasicGetAttributesResponse response = BasicGetAttributesResponse.onSuccess(request.getMsgType(),
request.getRequestId(), BasicAttributeKVMsg.from(result.get(0), result.get(1)));
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(response, src.getSessionId()), src.getServerAddress());
}
@Override
public void onFailure(Throwable t) {
if (t instanceof Exception) {
ToDeviceMsg toDeviceMsg = BasicStatusCodeResponse.onError(SessionMsgType.GET_ATTRIBUTES_REQUEST, request.getRequestId(), (Exception) t);
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(toDeviceMsg, src.getSessionId()), src.getServerAddress());
} else {
logger.error("[{}] Failed to process attributes request", deviceId, t);
}
}
});
}
private Optional<Set<String>> toOptionalSet(List<String> strings) {
if (strings == null || strings.isEmpty()) {
return Optional.empty();
} else {
return Optional.of(new HashSet<>(strings));
}
}
private void handleGetAttributesRequest(DeviceToDeviceActorMsg src) {
GetAttributesRequest request = (GetAttributesRequest) src.getPayload();
ListenableFuture<List<AttributeKvEntry>> clientAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.CLIENT_SCOPE, request.getClientAttributeNames());
ListenableFuture<List<AttributeKvEntry>> sharedAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.SHARED_SCOPE, request.getSharedAttributeNames());
Futures.addCallback(Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture)), new FutureCallback<List<List<AttributeKvEntry>>>() { Futures.addCallback(Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture)), new FutureCallback<List<List<AttributeKvEntry>>>() {
@Override @Override
public void onSuccess(@Nullable List<List<AttributeKvEntry>> result) { public void onSuccess(@Nullable List<List<AttributeKvEntry>> result) {
BasicGetAttributesResponse response = BasicGetAttributesResponse.onSuccess(request.getMsgType(), GetAttributeResponseMsg responseMsg = GetAttributeResponseMsg.newBuilder()
request.getRequestId(), BasicAttributeKVMsg.from(result.get(0), result.get(1))); .setRequestId(requestId)
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(response, src.getSessionId()), src.getServerAddress()); .addAllClientAttributeList(toTsKvProtos(result.get(0)))
.addAllSharedAttributeList(toTsKvProtos(result.get(1)))
.build();
sendToTransport(responseMsg, sessionId, sessionInfo);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
if (t instanceof Exception) { GetAttributeResponseMsg responseMsg = GetAttributeResponseMsg.newBuilder()
ToDeviceMsg toDeviceMsg = BasicStatusCodeResponse.onError(SessionMsgType.GET_ATTRIBUTES_REQUEST, request.getRequestId(), (Exception) t); .setError(t.getMessage())
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(toDeviceMsg, src.getSessionId()), src.getServerAddress()); .build();
} else { sendToTransport(responseMsg, sessionId, sessionInfo);
logger.error("[{}] Failed to process attributes request", deviceId, t);
}
} }
}); });
} }
@ -376,36 +314,36 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
} }
private void handleClientSideRPCRequest(ActorContext context, DeviceToDeviceActorMsg src) { // private void handleClientSideRPCRequest(ActorContext context, DeviceToDeviceActorMsg src) {
ToServerRpcRequestMsg request = (ToServerRpcRequestMsg) src.getPayload(); // ToServerRpcRequestMsg request = (ToServerRpcRequestMsg) src.getPayload();
//
JsonObject json = new JsonObject(); // JsonObject json = new JsonObject();
json.addProperty("method", request.getMethod()); // json.addProperty("method", request.getMethod());
json.add("params", jsonParser.parse(request.getParams())); // json.add("params", jsonParser.parse(request.getParams()));
//
TbMsgMetaData requestMetaData = defaultMetaData.copy(); // TbMsgMetaData requestMetaData = defaultMetaData.copy();
requestMetaData.putValue("requestId", Integer.toString(request.getRequestId())); // requestMetaData.putValue("requestId", Integer.toString(request.getRequestId()));
TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, requestMetaData, TbMsgDataType.JSON, gson.toJson(json), null, null, 0L); // TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, requestMetaData, TbMsgDataType.JSON, gson.toJson(json), null, null, 0L);
PendingSessionMsgData msgData = new PendingSessionMsgData(src.getSessionId(), src.getServerAddress(), SessionMsgType.TO_SERVER_RPC_REQUEST, request.getRequestId(), false, 1); // PendingSessionMsgData msgData = new PendingSessionMsgData(src.getSessionId(), src.getServerAddress(), SessionMsgType.TO_SERVER_RPC_REQUEST, request.getRequestId(), false, 1);
pushToRuleEngineWithTimeout(context, tbMsg, msgData); // pushToRuleEngineWithTimeout(context, tbMsg, msgData);
//
scheduleMsgWithDelay(context, new DeviceActorClientSideRpcTimeoutMsg(request.getRequestId(), systemContext.getClientSideRpcTimeout()), systemContext.getClientSideRpcTimeout()); // scheduleMsgWithDelay(context, new DeviceActorClientSideRpcTimeoutMsg(request.getRequestId(), systemContext.getClientSideRpcTimeout()), systemContext.getClientSideRpcTimeout());
toServerRpcPendingMap.put(request.getRequestId(), new ToServerRpcRequestMetadata(src.getSessionId(), src.getSessionType(), src.getServerAddress())); // toServerRpcPendingMap.put(request.getRequestId(), new ToServerRpcRequestMetadata(src.getSessionId(), src.getSessionType(), src.getServerAddress()));
} // }
public void processClientSideRpcTimeout(ActorContext context, DeviceActorClientSideRpcTimeoutMsg msg) { public void processClientSideRpcTimeout(ActorContext context, DeviceActorClientSideRpcTimeoutMsg msg) {
ToServerRpcRequestMetadata data = toServerRpcPendingMap.remove(msg.getId()); ToServerRpcRequestMetadata data = toServerRpcPendingMap.remove(msg.getId());
if (data != null) { if (data != null) {
logger.debug("[{}] Client side RPC request [{}] timeout detected!", deviceId, msg.getId()); logger.debug("[{}] Client side RPC request [{}] timeout detected!", deviceId, msg.getId());
ToDeviceMsg toDeviceMsg = new RuleEngineErrorMsg(SessionMsgType.TO_SERVER_RPC_REQUEST, RuleEngineError.TIMEOUT); ToDeviceMsg toDeviceMsg = new RuleEngineErrorMsg(SessionMsgType.TO_SERVER_RPC_REQUEST, RuleEngineError.TIMEOUT);
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(toDeviceMsg, data.getSessionId()), data.getServer()); // sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(toDeviceMsg, data.getSessionId()), data.getServer());
} }
} }
void processToServerRPCResponse(ActorContext context, ToServerRpcResponseActorMsg msg) { void processToServerRPCResponse(ActorContext context, ToServerRpcResponseActorMsg msg) {
ToServerRpcRequestMetadata data = toServerRpcPendingMap.remove(msg.getMsg().getRequestId()); ToServerRpcRequestMetadata data = toServerRpcPendingMap.remove(msg.getMsg().getRequestId());
if (data != null) { if (data != null) {
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(msg.getMsg(), data.getSessionId()), data.getServer()); // sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(msg.getMsg(), data.getSessionId()), data.getServer());
} }
} }
@ -433,68 +371,68 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
if (notification != null) { if (notification != null) {
ToDeviceMsg finalNotification = notification; ToDeviceMsg finalNotification = notification;
attributeSubscriptions.entrySet().forEach(sub -> { // attributeSubscriptions.entrySet().forEach(sub -> {
ActorSystemToDeviceSessionActorMsg response = new BasicActorSystemToDeviceSessionActorMsg(finalNotification, sub.getKey()); // ActorSystemToDeviceSessionActorMsg response = new BasicActorSystemToDeviceSessionActorMsg(finalNotification, sub.getKey());
sendMsgToSessionActor(response, sub.getValue().getServer()); // sendMsgToSessionActor(response, sub.getValue().getServer());
}); // });
} }
} else { } else {
logger.debug("[{}] No registered attributes subscriptions to process!", deviceId); logger.debug("[{}] No registered attributes subscriptions to process!", deviceId);
} }
} }
private void processRpcResponses(ActorContext context, DeviceToDeviceActorMsg msg) { // private void processRpcResponses(ActorContext context, DeviceToDeviceActorMsg msg) {
SessionId sessionId = msg.getSessionId(); // SessionId sessionId = msg.getSessionId();
FromDeviceMsg inMsg = msg.getPayload(); // FromDeviceMsg inMsg = msg.getPayload();
if (inMsg.getMsgType() == SessionMsgType.TO_DEVICE_RPC_RESPONSE) { // if (inMsg.getMsgType() == SessionMsgType.TO_DEVICE_RPC_RESPONSE) {
logger.debug("[{}] Processing rpc command response [{}]", deviceId, sessionId); // logger.debug("[{}] Processing rpc command response [{}]", deviceId, sessionId);
ToDeviceRpcResponseMsg responseMsg = (ToDeviceRpcResponseMsg) inMsg; // ToDeviceRpcResponseMsg responseMsg = (ToDeviceRpcResponseMsg) inMsg;
ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); // ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId());
boolean success = requestMd != null; // boolean success = requestMd != null;
if (success) { // if (success) {
systemContext.getDeviceRpcService().processRpcResponseFromDevice(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), // systemContext.getDeviceRpcService().processRpcResponseFromDevice(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(),
requestMd.getMsg().getServerAddress(), responseMsg.getData(), null)); // requestMd.getMsg().getServerAddress(), responseMsg.getData(), null));
} else { // } else {
logger.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); // logger.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId());
} // }
if (msg.getSessionType() == SessionType.SYNC) { // if (msg.getSessionType() == SessionType.SYNC) {
BasicCommandAckResponse response = success // BasicCommandAckResponse response = success
? BasicCommandAckResponse.onSuccess(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId()) // ? BasicCommandAckResponse.onSuccess(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId())
: BasicCommandAckResponse.onError(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId(), new TimeoutException()); // : BasicCommandAckResponse.onError(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId(), new TimeoutException());
sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(response, msg.getSessionId()), msg.getServerAddress()); // sendMsgToSessionActor(new BasicActorSystemToDeviceSessionActorMsg(response, msg.getSessionId()), msg.getServerAddress());
} // }
} // }
} // }
void processClusterEventMsg(ClusterEventMsg msg) { void processClusterEventMsg(ClusterEventMsg msg) {
if (!msg.isAdded()) { // if (!msg.isAdded()) {
logger.debug("[{}] Clearing attributes/rpc subscription for server [{}]", deviceId, msg.getServerAddress()); // logger.debug("[{}] Clearing attributes/rpc subscription for server [{}]", deviceId, msg.getServerAddress());
Predicate<Map.Entry<SessionId, SessionInfo>> filter = e -> e.getValue().getServer() // Predicate<Map.Entry<SessionId, SessionInfo>> filter = e -> e.getValue().getServer()
.map(serverAddress -> serverAddress.equals(msg.getServerAddress())).orElse(false); // .map(serverAddress -> serverAddress.equals(msg.getServerAddress())).orElse(false);
attributeSubscriptions.entrySet().removeIf(filter); // attributeSubscriptions.entrySet().removeIf(filter);
rpcSubscriptions.entrySet().removeIf(filter); // rpcSubscriptions.entrySet().removeIf(filter);
} // }
} }
private void processSubscriptionCommands(ActorContext context, DeviceToDeviceActorMsg msg) { // private void processSubscriptionCommands(ActorContext context, DeviceToDeviceActorMsg msg) {
SessionId sessionId = msg.getSessionId(); // SessionId sessionId = msg.getSessionId();
SessionType sessionType = msg.getSessionType(); // SessionType sessionType = msg.getSessionType();
FromDeviceMsg inMsg = msg.getPayload(); // FromDeviceMsg inMsg = msg.getPayload();
if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_ATTRIBUTES_REQUEST) { // if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_ATTRIBUTES_REQUEST) {
logger.debug("[{}] Registering attributes subscription for session [{}]", deviceId, sessionId); // logger.debug("[{}] Registering attributes subscription for session [{}]", deviceId, sessionId);
attributeSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress())); // attributeSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress()));
} else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_ATTRIBUTES_REQUEST) { // } else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_ATTRIBUTES_REQUEST) {
logger.debug("[{}] Canceling attributes subscription for session [{}]", deviceId, sessionId); // logger.debug("[{}] Canceling attributes subscription for session [{}]", deviceId, sessionId);
attributeSubscriptions.remove(sessionId); // attributeSubscriptions.remove(sessionId);
} else if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_RPC_COMMANDS_REQUEST) { // } else if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_RPC_COMMANDS_REQUEST) {
logger.debug("[{}] Registering rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType); // logger.debug("[{}] Registering rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType);
rpcSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress())); // rpcSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress()));
sendPendingRequests(context, sessionId, sessionType, msg.getServerAddress()); // sendPendingRequests(context, sessionId, sessionType, msg.getServerAddress());
} else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_RPC_COMMANDS_REQUEST) { // } else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_RPC_COMMANDS_REQUEST) {
logger.debug("[{}] Canceling rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType); // logger.debug("[{}] Canceling rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType);
rpcSubscriptions.remove(sessionId); // rpcSubscriptions.remove(sessionId);
} // }
} // }
private void processSessionStateMsgs(SessionInfoProto sessionInfo, SessionEventMsg msg) { private void processSessionStateMsgs(SessionInfoProto sessionInfo, SessionEventMsg msg) {
UUID sessionId = new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()); UUID sessionId = new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB());
@ -506,15 +444,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
closeSession(sessionIdToRemove, sessions.remove(sessionIdToRemove)); closeSession(sessionIdToRemove, sessions.remove(sessionIdToRemove));
} }
} }
sessions.put(sessionId, new SessionInfo(SessionType.ASYNC, msg.getServerAddress())); sessions.put(sessionId, new SessionInfo(TransportProtos.SessionType.ASYNC, sessionInfo.getNodeId()));
if (sessions.size() == 1) { if (sessions.size() == 1) {
reportSessionOpen(); reportSessionOpen();
} }
} } else if (msg.getEvent() == SessionEvent.CLOSED) {
FromDeviceMsg inMsg = msg.getPayload();
if (inMsg instanceof SessionOpenMsg) {
logger.debug("[{}] Processing new session [{}]", deviceId, sessionId);
} else if (inMsg instanceof SessionCloseMsg) {
logger.debug("[{}] Canceling subscriptions for closed session [{}]", deviceId, sessionId); logger.debug("[{}] Canceling subscriptions for closed session [{}]", deviceId, sessionId);
sessions.remove(sessionId); sessions.remove(sessionId);
attributeSubscriptions.remove(sessionId); attributeSubscriptions.remove(sessionId);
@ -532,7 +466,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
systemContext.getRpcService().tell(systemContext.getEncodingService() systemContext.getRpcService().tell(systemContext.getEncodingService()
.convertToProtoDataMessage(sessionAddress.get(), response)); .convertToProtoDataMessage(sessionAddress.get(), response));
} else { } else {
systemContext.getSessionManagerActor().tell(response, ActorRef.noSender()); // systemContext.getSessionManagerActor().tell(response, ActorRef.noSender());
} }
} }
@ -578,4 +512,62 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
return json; return json;
} }
private Optional<Set<String>> toOptionalSet(List<String> strings) {
if (strings == null || strings.isEmpty()) {
return Optional.empty();
} else {
return Optional.of(new HashSet<>(strings));
}
}
private void sendToTransport(GetAttributeResponseMsg responseMsg, UUID sessionId, SessionInfoProto sessionInfo) {
DeviceActorToTransportMsg msg = DeviceActorToTransportMsg.newBuilder()
.setSessionIdMSB(sessionId.getMostSignificantBits())
.setSessionIdLSB(sessionId.getLeastSignificantBits())
.setGetAttributesResponse(responseMsg).build();
systemContext.getRuleEngineTransportService().process(sessionInfo.getNodeId(), msg);
}
private List<TsKvProto> toTsKvProtos(@Nullable List<AttributeKvEntry> result) {
List<TsKvProto> clientAttributes;
if (result == null || result.isEmpty()) {
clientAttributes = Collections.emptyList();
} else {
clientAttributes = new ArrayList<>(result.size());
for (AttributeKvEntry attrEntry : result) {
clientAttributes.add(toTsKvProto(attrEntry));
}
}
return clientAttributes;
}
private TsKvProto toTsKvProto(AttributeKvEntry attrEntry) {
return TsKvProto.newBuilder().setTs(attrEntry.getLastUpdateTs())
.setKv(toKeyValueProto(attrEntry)).build();
}
private KeyValueProto toKeyValueProto(KvEntry kvEntry) {
KeyValueProto.Builder builder = KeyValueProto.newBuilder();
builder.setKey(kvEntry.getKey());
switch (kvEntry.getDataType()) {
case BOOLEAN:
builder.setType(KeyValueType.BOOLEAN_V);
builder.setBoolV(kvEntry.getBooleanValue().get());
break;
case DOUBLE:
builder.setType(KeyValueType.DOUBLE_V);
builder.setDoubleV(kvEntry.getDoubleValue().get());
break;
case LONG:
builder.setType(KeyValueType.LONG_V);
builder.setLongV(kvEntry.getLongValue().get());
break;
case STRING:
builder.setType(KeyValueType.STRING_V);
builder.setStringV(kvEntry.getStrValue().get());
break;
}
return builder.build();
}
} }

36
application/src/main/java/org/thingsboard/server/actors/device/RuleEngineQueuePutAckMsg.java

@ -1,36 +0,0 @@
/**
* 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.actors.device;
import lombok.Data;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg;
import java.util.UUID;
/**
* Created by ashvayka on 15.03.18.
*/
@Data
public final class RuleEngineQueuePutAckMsg implements TbActorMsg {
private final UUID id;
@Override
public MsgType getMsgType() {
return MsgType.RULE_ENGINE_QUEUE_PUT_ACK_MSG;
}
}

8
application/src/main/java/org/thingsboard/server/actors/device/SessionInfo.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

40
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java

@ -25,7 +25,6 @@ import java.util.Optional;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.device.DeviceActorToRuleEngineMsg; import org.thingsboard.server.actors.device.DeviceActorToRuleEngineMsg;
import org.thingsboard.server.actors.device.RuleEngineQueuePutAckMsg;
import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.actors.shared.ComponentMsgProcessor; import org.thingsboard.server.actors.shared.ComponentMsgProcessor;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
@ -90,26 +89,12 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode)); nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode));
} }
initRoutes(ruleChain, ruleNodeList); initRoutes(ruleChain, ruleNodeList);
reprocess(ruleNodeList);
started = true; started = true;
} else { } else {
onUpdate(context); onUpdate(context);
} }
} }
private void reprocess(List<RuleNode> ruleNodeList) {
for (RuleNode ruleNode : ruleNodeList) {
for (TbMsg tbMsg : queue.findUnprocessed(tenantId, ruleNode.getId().getId(), systemContext.getQueuePartitionId())) {
pushMsgToNode(nodeActors.get(ruleNode.getId()), tbMsg, "");
}
}
if (firstNode != null) {
for (TbMsg tbMsg : queue.findUnprocessed(tenantId, entityId.getId(), systemContext.getQueuePartitionId())) {
pushMsgToNode(firstNode, tbMsg, "");
}
}
}
@Override @Override
public void onUpdate(ActorContext context) throws Exception { public void onUpdate(ActorContext context) throws Exception {
RuleChain ruleChain = service.findRuleChainById(entityId); RuleChain ruleChain = service.findRuleChainById(entityId);
@ -134,7 +119,6 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
}); });
initRoutes(ruleChain, ruleNodeList); initRoutes(ruleChain, ruleNodeList);
reprocess(ruleNodeList);
} }
@Override @Override
@ -188,17 +172,14 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
void onServiceToRuleEngineMsg(ServiceToRuleEngineMsg envelope) { void onServiceToRuleEngineMsg(ServiceToRuleEngineMsg envelope) {
checkActive(); checkActive();
if (firstNode != null) { if (firstNode != null) {
putToQueue(enrichWithRuleChainId(envelope.getTbMsg()), msg -> pushMsgToNode(firstNode, msg, "")); pushMsgToNode(firstNode, enrichWithRuleChainId(envelope.getTbMsg()), "");
} }
} }
void onDeviceActorToRuleEngineMsg(DeviceActorToRuleEngineMsg envelope) { void onDeviceActorToRuleEngineMsg(DeviceActorToRuleEngineMsg envelope) {
checkActive(); checkActive();
if (firstNode != null) { if (firstNode != null) {
putToQueue(enrichWithRuleChainId(envelope.getTbMsg()), msg -> { pushMsgToNode(firstNode, enrichWithRuleChainId(envelope.getTbMsg()), "");
pushMsgToNode(firstNode, msg, "");
envelope.getCallbackRef().tell(new RuleEngineQueuePutAckMsg(msg.getId()), self);
});
} }
} }
@ -206,15 +187,16 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
checkActive(); checkActive();
if (envelope.isEnqueue()) { if (envelope.isEnqueue()) {
if (firstNode != null) { if (firstNode != null) {
putToQueue(enrichWithRuleChainId(envelope.getMsg()), msg -> pushMsgToNode(firstNode, msg, envelope.getFromRelationType())); pushMsgToNode(firstNode, enrichWithRuleChainId(envelope.getMsg()), envelope.getFromRelationType());
} }
} else { } else {
if (firstNode != null) { if (firstNode != null) {
pushMsgToNode(firstNode, envelope.getMsg(), envelope.getFromRelationType()); pushMsgToNode(firstNode, envelope.getMsg(), envelope.getFromRelationType());
} else { } else {
TbMsg msg = envelope.getMsg(); // TODO: Ack this message in Kafka
EntityId ackId = msg.getRuleNodeId() != null ? msg.getRuleNodeId() : msg.getRuleChainId(); // TbMsg msg = envelope.getMsg();
queue.ack(tenantId, envelope.getMsg(), ackId.getId(), msg.getClusterPartition()); // EntityId ackId = msg.getRuleNodeId() != null ? msg.getRuleNodeId() : msg.getRuleChainId();
// queue.ack(tenantId, envelope.getMsg(), ackId.getId(), msg.getClusterPartition());
} }
} }
} }
@ -249,7 +231,8 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
EntityId ackId = msg.getRuleNodeId() != null ? msg.getRuleNodeId() : msg.getRuleChainId(); EntityId ackId = msg.getRuleNodeId() != null ? msg.getRuleNodeId() : msg.getRuleChainId();
if (relationsCount == 0) { if (relationsCount == 0) {
if (ackId != null) { if (ackId != null) {
queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition()); // TODO: Ack this message in Kafka
// queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition());
} }
} else if (relationsCount == 1) { } else if (relationsCount == 1) {
for (RuleNodeRelation relation : relations) { for (RuleNodeRelation relation : relations) {
@ -269,7 +252,8 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
} }
//TODO: Ideally this should happen in async way when all targets confirm that the copied messages are successfully written to corresponding target queues. //TODO: Ideally this should happen in async way when all targets confirm that the copied messages are successfully written to corresponding target queues.
if (ackId != null) { if (ackId != null) {
queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition()); // TODO: Ack this message in Kafka
// queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition());
} }
} }
} }
@ -296,7 +280,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
RuleNodeId targetId = new RuleNodeId(target.getId()); RuleNodeId targetId = new RuleNodeId(target.getId());
RuleNodeCtx targetNodeCtx = nodeActors.get(targetId); RuleNodeCtx targetNodeCtx = nodeActors.get(targetId);
TbMsg copy = msg.copy(UUIDs.timeBased(), entityId, targetId, DEFAULT_CLUSTER_PARTITION); TbMsg copy = msg.copy(UUIDs.timeBased(), entityId, targetId, DEFAULT_CLUSTER_PARTITION);
putToQueue(copy, queuedMsg -> pushMsgToNode(targetNodeCtx, queuedMsg, fromRelationType)); pushMsgToNode(targetNodeCtx, copy, fromRelationType);
} }
private void pushToTarget(TbMsg msg, EntityId target, String fromRelationType) { private void pushToTarget(TbMsg msg, EntityId target, String fromRelationType) {

14
application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java

@ -30,7 +30,6 @@ import org.thingsboard.server.actors.app.AppActor;
import org.thingsboard.server.actors.rpc.RpcBroadcastMsg; import org.thingsboard.server.actors.rpc.RpcBroadcastMsg;
import org.thingsboard.server.actors.rpc.RpcManagerActor; import org.thingsboard.server.actors.rpc.RpcManagerActor;
import org.thingsboard.server.actors.rpc.RpcSessionCreateRequestMsg; import org.thingsboard.server.actors.rpc.RpcSessionCreateRequestMsg;
import org.thingsboard.server.actors.session.SessionManagerActor;
import org.thingsboard.server.actors.stats.StatsActor; import org.thingsboard.server.actors.stats.StatsActor;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
@ -90,8 +89,6 @@ public class DefaultActorService implements ActorService {
private ActorRef appActor; private ActorRef appActor;
private ActorRef sessionManagerActor;
private ActorRef rpcManagerActor; private ActorRef rpcManagerActor;
@PostConstruct @PostConstruct
@ -104,10 +101,6 @@ public class DefaultActorService implements ActorService {
appActor = system.actorOf(Props.create(new AppActor.ActorCreator(actorContext)).withDispatcher(APP_DISPATCHER_NAME), "appActor"); appActor = system.actorOf(Props.create(new AppActor.ActorCreator(actorContext)).withDispatcher(APP_DISPATCHER_NAME), "appActor");
actorContext.setAppActor(appActor); actorContext.setAppActor(appActor);
sessionManagerActor = system.actorOf(Props.create(new SessionManagerActor.ActorCreator(actorContext)).withDispatcher(CORE_DISPATCHER_NAME),
"sessionManagerActor");
actorContext.setSessionManagerActor(sessionManagerActor);
rpcManagerActor = system.actorOf(Props.create(new RpcManagerActor.ActorCreator(actorContext)).withDispatcher(CORE_DISPATCHER_NAME), rpcManagerActor = system.actorOf(Props.create(new RpcManagerActor.ActorCreator(actorContext)).withDispatcher(CORE_DISPATCHER_NAME),
"rpcManagerActor"); "rpcManagerActor");
@ -134,12 +127,6 @@ public class DefaultActorService implements ActorService {
appActor.tell(msg, ActorRef.noSender()); appActor.tell(msg, ActorRef.noSender());
} }
@Override
public void process(SessionAwareMsg msg) {
log.debug("Processing session aware msg: {}", msg);
sessionManagerActor.tell(msg, ActorRef.noSender());
}
@Override @Override
public void onServerAdded(ServerInstance server) { public void onServerAdded(ServerInstance server) {
log.trace("Processing onServerAdded msg: {}", server); log.trace("Processing onServerAdded msg: {}", server);
@ -194,7 +181,6 @@ public class DefaultActorService implements ActorService {
private void broadcast(ClusterEventMsg msg) { private void broadcast(ClusterEventMsg msg) {
this.appActor.tell(msg, ActorRef.noSender()); this.appActor.tell(msg, ActorRef.noSender());
this.sessionManagerActor.tell(msg, ActorRef.noSender());
this.rpcManagerActor.tell(msg, ActorRef.noSender()); this.rpcManagerActor.tell(msg, ActorRef.noSender());
} }

16
application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java

@ -35,14 +35,12 @@ public abstract class ComponentMsgProcessor<T extends EntityId> extends Abstract
protected final TenantId tenantId; protected final TenantId tenantId;
protected final T entityId; protected final T entityId;
protected final MsgQueueService queue;
protected ComponentLifecycleState state; protected ComponentLifecycleState state;
protected ComponentMsgProcessor(ActorSystemContext systemContext, LoggingAdapter logger, TenantId tenantId, T id) { protected ComponentMsgProcessor(ActorSystemContext systemContext, LoggingAdapter logger, TenantId tenantId, T id) {
super(systemContext, logger); super(systemContext, logger);
this.tenantId = tenantId; this.tenantId = tenantId;
this.entityId = id; this.entityId = id;
this.queue = systemContext.getMsgQueueService();
} }
public abstract void start(ActorContext context) throws Exception; public abstract void start(ActorContext context) throws Exception;
@ -86,18 +84,4 @@ public abstract class ComponentMsgProcessor<T extends EntityId> extends Abstract
} }
} }
protected void putToQueue(final TbMsg tbMsg, final Consumer<TbMsg> onSuccess) {
EntityId entityId = tbMsg.getRuleNodeId() != null ? tbMsg.getRuleNodeId() : tbMsg.getRuleChainId();
Futures.addCallback(queue.put(this.tenantId, tbMsg, entityId.getId(), tbMsg.getClusterPartition()), new FutureCallback<Void>() {
@Override
public void onSuccess(@Nullable Void result) {
onSuccess.accept(tbMsg);
}
@Override
public void onFailure(Throwable t) {
logger.debug("Failed to push message [{}] to queue due to [{}]", tbMsg, t);
}
});
}
} }

17
application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java

@ -1,3 +1,18 @@
/**
* 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.service.transport; package org.thingsboard.server.service.transport;
import akka.actor.ActorRef; import akka.actor.ActorRef;
@ -30,6 +45,7 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.time.Duration; import java.time.Duration;
import java.util.Optional; import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.function.Consumer; import java.util.function.Consumer;
@ -136,6 +152,7 @@ public class RemoteRuleEngineTransportService implements RuleEngineTransportServ
@Override @Override
public void process(String nodeId, DeviceActorToTransportMsg msg, Runnable onSuccess, Consumer<Throwable> onFailure) { public void process(String nodeId, DeviceActorToTransportMsg msg, Runnable onSuccess, Consumer<Throwable> onFailure) {
notificationsProducer.send(notificationsTopic + "." + nodeId, notificationsProducer.send(notificationsTopic + "." + nodeId,
new UUID(msg.getSessionIdMSB(), msg.getSessionIdLSB()).toString(),
ToTransportMsg.newBuilder().setToDeviceSessionMsg(msg).build() ToTransportMsg.newBuilder().setToDeviceSessionMsg(msg).build()
, new QueueCallbackAdaptor(onSuccess, onFailure)); , new QueueCallbackAdaptor(onSuccess, onFailure));
} }

9
application/src/main/java/org/thingsboard/server/service/transport/RuleEngineTransportService.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -16,7 +16,6 @@
package org.thingsboard.server.service.transport; package org.thingsboard.server.service.transport;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceActorToTransportMsg; import org.thingsboard.server.gen.transport.TransportProtos.DeviceActorToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.;
import java.util.function.Consumer; import java.util.function.Consumer;

8
application/src/main/java/org/thingsboard/server/service/transport/ToTransportMsgEncoder.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

16
application/src/main/java/org/thingsboard/server/service/transport/msg/TransportToDeviceActorMsgWrapper.java

@ -1,3 +1,18 @@
/**
* 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.service.transport.msg; package org.thingsboard.server.service.transport.msg;
import lombok.Data; import lombok.Data;
@ -7,7 +22,6 @@ import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg;
import org.thingsboard.server.common.msg.aware.TenantAwareMsg; import org.thingsboard.server.common.msg.aware.TenantAwareMsg;
import org.thingsboard.server.common.msg.cluster.ServerAddress;
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg;
import java.io.Serializable; import java.io.Serializable;

3
common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java

@ -91,8 +91,6 @@ public enum MsgType {
DEVICE_ACTOR_CLIENT_SIDE_RPC_TIMEOUT_MSG, DEVICE_ACTOR_CLIENT_SIDE_RPC_TIMEOUT_MSG,
DEVICE_ACTOR_QUEUE_TIMEOUT_MSG,
/** /**
* Message that is sent from the Device Actor to Rule Engine. Requires acknowledgement * Message that is sent from the Device Actor to Rule Engine. Requires acknowledgement
*/ */
@ -101,7 +99,6 @@ public enum MsgType {
/** /**
* Message that is sent from Rule Engine to the Device Actor when message is successfully pushed to queue. * Message that is sent from Rule Engine to the Device Actor when message is successfully pushed to queue.
*/ */
RULE_ENGINE_QUEUE_PUT_ACK_MSG,
ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG, ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG,
TRANSPORT_TO_DEVICE_SESSION_ACTOR_MSG, TRANSPORT_TO_DEVICE_SESSION_ACTOR_MSG,
SESSION_TIMEOUT_MSG, SESSION_TIMEOUT_MSG,

52
common/message/src/main/java/org/thingsboard/server/common/msg/core/BasicActorSystemToDeviceSessionActorMsg.java

@ -1,52 +0,0 @@
/**
* 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.common.msg.core;
import org.thingsboard.server.common.data.id.SessionId;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.session.ToDeviceMsg;
public class BasicActorSystemToDeviceSessionActorMsg implements ActorSystemToDeviceSessionActorMsg {
private final ToDeviceMsg msg;
private final SessionId sessionId;
public BasicActorSystemToDeviceSessionActorMsg(ToDeviceMsg msg, SessionId sessionId) {
super();
this.msg = msg;
this.sessionId = sessionId;
}
@Override
public SessionId getSessionId() {
return sessionId;
}
@Override
public ToDeviceMsg getMsg() {
return msg;
}
@Override
public String toString() {
return "BasicActorSystemToDeviceSessionActorMsg [msg=" + msg + ", sessionId=" + sessionId + "]";
}
@Override
public MsgType getMsgType() {
return MsgType.ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG;
}
}

36
common/message/src/main/java/org/thingsboard/server/common/msg/timeout/DeviceActorQueueTimeoutMsg.java

@ -1,36 +0,0 @@
/**
* 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.common.msg.timeout;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.timeout.TimeoutMsg;
import java.util.UUID;
/**
* @author Andrew Shvayka
*/
public final class DeviceActorQueueTimeoutMsg extends TimeoutMsg<UUID> {
public DeviceActorQueueTimeoutMsg(UUID id, long timeout) {
super(id, timeout);
}
@Override
public MsgType getMsgType() {
return MsgType.DEVICE_ACTOR_QUEUE_TIMEOUT_MSG;
}
}

8
common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaConsumerTemplate.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

12
common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -106,6 +106,10 @@ public class TBKafkaProducerTemplate<T> {
return send(topic, key, value, null, headers, callback); return send(topic, key, value, null, headers, callback);
} }
public Future<RecordMetadata> send(String topic, String key, T value, Callback callback) {
return send(topic, key, value, null, null, callback);
}
public Future<RecordMetadata> send(String topic, String key, T value, Long timestamp, Iterable<Header> headers, Callback callback) { public Future<RecordMetadata> send(String topic, String key, T value, Long timestamp, Iterable<Header> headers, Callback callback) {
byte[] data = encoder.encode(value); byte[] data = encoder.encode(value);
ProducerRecord<String, byte[]> record; ProducerRecord<String, byte[]> record;

2
common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java

@ -160,7 +160,7 @@ public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTe
SettableFuture<Response> future = SettableFuture.create(); SettableFuture<Response> future = SettableFuture.create();
pendingRequests.putIfAbsent(requestId, new ResponseMetaData<>(tickTs + maxRequestTimeout, future)); pendingRequests.putIfAbsent(requestId, new ResponseMetaData<>(tickTs + maxRequestTimeout, future));
request = requestTemplate.enrich(request, responseTemplate.getTopic(), requestId); request = requestTemplate.enrich(request, responseTemplate.getTopic(), requestId);
requestTemplate.send(key, request, headers); requestTemplate.send(key, request, headers, null);
return future; return future;
} }

8
common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

8
common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

2
common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgProcessor.java

@ -20,8 +20,6 @@ import org.thingsboard.server.common.msg.aware.SessionAwareMsg;
public interface SessionMsgProcessor { public interface SessionMsgProcessor {
void process(SessionAwareMsg msg);
void onDeviceAdded(Device device); void onDeviceAdded(Device device);
} }

1
common/transport/src/main/java/org/thingsboard/server/common/transport/TransportAdaptor.java

@ -20,6 +20,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg;
import org.thingsboard.server.common.msg.session.SessionContext; import org.thingsboard.server.common.msg.session.SessionContext;
import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Optional; import java.util.Optional;

11
common/transport/src/main/java/org/thingsboard/server/common/transport/TransportService.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -22,6 +22,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionEventMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
/** /**
* Created by ashvayka on 04.10.18. * Created by ashvayka on 04.10.18.
@ -40,6 +41,8 @@ public interface TransportService {
void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TransportServiceCallback<Void> callback); void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback);
void registerSession(SessionInfoProto sessionInfo, SessionMsgListener listener); void registerSession(SessionInfoProto sessionInfo, SessionMsgListener listener);
void deregisterSession(SessionInfoProto sessionInfo); void deregisterSession(SessionInfoProto sessionInfo);

180
common/transport/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java

@ -15,18 +15,38 @@
*/ */
package org.thingsboard.server.common.transport.adaptor; package org.thingsboard.server.common.transport.adaptor;
import com.google.gson.Gson;
import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.JsonPrimitive;
import com.google.gson.JsonSyntaxException;
import org.thingsboard.server.common.data.kv.AttributeKey;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.msg.core.AttributesUpdateRequest;
import org.thingsboard.server.common.msg.core.BasicAttributesUpdateRequest;
import org.thingsboard.server.common.msg.core.BasicRequest;
import org.thingsboard.server.common.msg.core.BasicTelemetryUploadRequest;
import org.thingsboard.server.common.msg.core.TelemetryUploadRequest;
import org.thingsboard.server.common.msg.core.ToDeviceRpcRequestMsg;
import org.thingsboard.server.common.msg.core.ToServerRpcRequestMsg;
import org.thingsboard.server.common.msg.core.ToServerRpcResponseMsg;
import org.thingsboard.server.common.msg.kv.AttributesKVMsg;
import org.thingsboard.server.gen.transport.TransportProtos.*;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map.Entry; import java.util.Map.Entry;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import com.google.gson.*;
import org.thingsboard.server.common.msg.core.*;
import org.thingsboard.server.common.data.kv.*;
import org.thingsboard.server.common.msg.kv.AttributesKVMsg;
public class JsonConverter { public class JsonConverter {
private static final Gson GSON = new Gson(); private static final Gson GSON = new Gson();
@ -44,6 +64,109 @@ public class JsonConverter {
return convertToTelemetry(jsonObject, System.currentTimeMillis(), requestId); return convertToTelemetry(jsonObject, System.currentTimeMillis(), requestId);
} }
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonObject) throws JsonSyntaxException {
long systemTs = System.currentTimeMillis();
PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder();
if (jsonObject.isJsonObject()) {
parseObject(builder, systemTs, jsonObject);
} else if (jsonObject.isJsonArray()) {
jsonObject.getAsJsonArray().forEach(je -> {
if (je.isJsonObject()) {
parseObject(builder, systemTs, je.getAsJsonObject());
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + je);
}
});
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + jsonObject);
}
return builder.build();
}
public static PostAttributeMsg convertToAttributesProto(JsonElement jsonObject) throws JsonSyntaxException {
if (jsonObject.isJsonObject()) {
PostAttributeMsg.Builder result = PostAttributeMsg.newBuilder();
List<KeyValueProto> keyValueList = parseProtoValues(jsonObject.getAsJsonObject());
result.addAllKv(keyValueList);
return result.build();
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + jsonObject);
}
}
private static void parseObject(PostTelemetryMsg.Builder builder, long systemTs, JsonElement jsonObject) {
JsonObject jo = jsonObject.getAsJsonObject();
if (jo.has("ts") && jo.has("values")) {
parseWithTs(builder, jo);
} else {
parseWithoutTs(builder, systemTs, jo);
}
}
private static void parseWithoutTs(PostTelemetryMsg.Builder request, long systemTs, JsonObject jo) {
TsKvListProto.Builder builder = TsKvListProto.newBuilder();
builder.setTs(systemTs);
builder.addAllKv(parseProtoValues(jo));
request.addTsKvList(builder.build());
}
public static void parseWithTs(PostTelemetryMsg.Builder request, JsonObject jo) {
TsKvListProto.Builder builder = TsKvListProto.newBuilder();
builder.setTs(jo.get("ts").getAsLong());
builder.addAllKv(parseProtoValues(jo.get("values").getAsJsonObject()));
request.addTsKvList(builder.build());
}
public static List<KeyValueProto> parseProtoValues(JsonObject valuesObject) {
List<KeyValueProto> result = new ArrayList<>();
for (Entry<String, JsonElement> valueEntry : valuesObject.entrySet()) {
JsonElement element = valueEntry.getValue();
if (element.isJsonPrimitive()) {
JsonPrimitive value = element.getAsJsonPrimitive();
if (value.isString()) {
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.STRING_V)
.setStringV(value.getAsString()).build());
} else if (value.isBoolean()) {
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.BOOLEAN_V)
.setBoolV(value.getAsBoolean()).build());
} else if (value.isNumber()) {
if (value.getAsString().contains(".")) {
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.DOUBLE_V)
.setDoubleV(value.getAsDouble()).build());
} else {
try {
long longValue = Long.parseLong(value.getAsString());
result.add(KeyValueProto.newBuilder().setKey(valueEntry.getKey()).setType(KeyValueType.LONG_V)
.setLongV(longValue).build());
} catch (NumberFormatException e) {
throw new JsonSyntaxException("Big integer values are not supported!");
}
}
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + value);
}
} else {
throw new JsonSyntaxException(CAN_T_PARSE_VALUE + element);
}
}
return result;
}
private static void parseNumericProto(List<KvEntry> result, Entry<String, JsonElement> valueEntry, JsonPrimitive value) {
if (value.getAsString().contains(".")) {
result.add(new DoubleDataEntry(valueEntry.getKey(), value.getAsDouble()));
} else {
try {
long longValue = Long.parseLong(value.getAsString());
result.add(new LongDataEntry(valueEntry.getKey(), longValue));
} catch (NumberFormatException e) {
throw new JsonSyntaxException("Big integer values are not supported!");
}
}
}
private static TelemetryUploadRequest convertToTelemetry(JsonElement jsonObject, long systemTs, int requestId) throws JsonSyntaxException { private static TelemetryUploadRequest convertToTelemetry(JsonElement jsonObject, long systemTs, int requestId) throws JsonSyntaxException {
BasicTelemetryUploadRequest request = new BasicTelemetryUploadRequest(requestId); BasicTelemetryUploadRequest request = new BasicTelemetryUploadRequest(requestId);
if (jsonObject.isJsonObject()) { if (jsonObject.isJsonObject()) {
@ -140,6 +263,26 @@ public class JsonConverter {
} }
} }
public static JsonObject toJson(GetAttributeResponseMsg payload) {
JsonObject result = new JsonObject();
if (payload.getClientAttributeListCount() > 0) {
JsonObject attrObject = new JsonObject();
payload.getClientAttributeListList().forEach(addToObjectFromProto(attrObject));
result.add("client", attrObject);
}
if (payload.getSharedAttributeListCount() > 0) {
JsonObject attrObject = new JsonObject();
payload.getSharedAttributeListList().forEach(addToObjectFromProto(attrObject));
result.add("shared", attrObject);
}
if (payload.getDeletedAttributeKeysCount() > 0) {
JsonArray attrObject = new JsonArray();
payload.getDeletedAttributeKeysList().forEach(attrObject::add);
result.add("deleted", attrObject);
}
return result;
}
public static JsonObject toJson(AttributesKVMsg payload, boolean asMap) { public static JsonObject toJson(AttributesKVMsg payload, boolean asMap) {
JsonObject result = new JsonObject(); JsonObject result = new JsonObject();
if (asMap) { if (asMap) {
@ -166,8 +309,29 @@ public class JsonConverter {
} }
private static Consumer<AttributeKey> addToObject(JsonArray result) { private static Consumer<AttributeKey> addToObject(JsonArray result) {
return key -> { return key -> result.add(key.getAttributeKey());
result.add(key.getAttributeKey()); }
private static Consumer<TsKvProto> addToObjectFromProto(JsonObject result) {
return de -> {
JsonPrimitive value;
switch (de.getKv().getType()) {
case BOOLEAN_V:
value = new JsonPrimitive(de.getKv().getBoolV());
break;
case DOUBLE_V:
value = new JsonPrimitive(de.getKv().getDoubleV());
break;
case LONG_V:
value = new JsonPrimitive(de.getKv().getLongV());
break;
case STRING_V:
value = new JsonPrimitive(de.getKv().getStringV());
break;
default:
throw new IllegalArgumentException("Unsupported data type: " + de.getKv().getType());
}
result.add(de.getKv().getKey(), value);
}; };
} }

8
common/transport/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

36
common/transport/src/main/proto/transport.proto

@ -23,12 +23,13 @@ option java_outer_classname = "TransportProtos";
* Data Structures; * Data Structures;
*/ */
message SessionInfoProto { message SessionInfoProto {
int64 sessionIdMSB = 1; string nodeId = 1;
int64 sessionIdLSB = 2; int64 sessionIdMSB = 2;
int64 tenantIdMSB = 3; int64 sessionIdLSB = 3;
int64 tenantIdLSB = 4; int64 tenantIdMSB = 4;
int64 deviceIdMSB = 5; int64 tenantIdLSB = 5;
int64 deviceIdLSB = 6; int64 deviceIdMSB = 6;
int64 deviceIdLSB = 7;
} }
enum SessionEvent { enum SessionEvent {
@ -57,6 +58,11 @@ message KeyValueProto {
string string_v = 6; string string_v = 6;
} }
message TsKvProto {
int64 ts = 1;
KeyValueProto kv = 2;
}
message TsKvListProto { message TsKvListProto {
int64 ts = 1; int64 ts = 1;
repeated KeyValueProto kv = 2; repeated KeyValueProto kv = 2;
@ -76,9 +82,8 @@ message DeviceInfoProto {
* Messages that use Data Structures; * Messages that use Data Structures;
*/ */
message SessionEventMsg { message SessionEventMsg {
string nodeId = 1; SessionType sessionType = 1;
SessionType sessionType = 2; SessionEvent event = 2;
SessionEvent event = 3;
} }
message PostTelemetryMsg { message PostTelemetryMsg {
@ -90,14 +95,17 @@ message PostAttributeMsg {
} }
message GetAttributeRequestMsg { message GetAttributeRequestMsg {
repeated string clientAttributeNames = 1; int32 requestId = 1;
repeated string sharedAttributeNames = 2; repeated string clientAttributeNames = 2;
repeated string sharedAttributeNames = 3;
} }
message GetAttributeResponseMsg { message GetAttributeResponseMsg {
repeated TsKvListProto clientAttributeList = 1; int32 requestId = 1;
repeated TsKvListProto sharedAttributeList = 2; repeated TsKvProto clientAttributeList = 2;
repeated string deletedAttributeKeys = 3; repeated TsKvProto sharedAttributeList = 3;
repeated string deletedAttributeKeys = 4;
string error = 5;
} }
message ValidateDeviceTokenRequestMsg { message ValidateDeviceTokenRequestMsg {

48
transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java

@ -106,30 +106,30 @@ public class CoapServerTest {
public static SessionMsgProcessor sessionMsgProcessor() { public static SessionMsgProcessor sessionMsgProcessor() {
return new SessionMsgProcessor() { return new SessionMsgProcessor() {
@Override // @Override
public void process(SessionAwareMsg toActorMsg) { // public void process(SessionAwareMsg toActorMsg) {
if (toActorMsg instanceof TransportToDeviceSessionActorMsg) { // if (toActorMsg instanceof TransportToDeviceSessionActorMsg) {
AdaptorToSessionActorMsg sessionMsg = ((TransportToDeviceSessionActorMsg) toActorMsg).getSessionMsg(); // AdaptorToSessionActorMsg sessionMsg = ((TransportToDeviceSessionActorMsg) toActorMsg).getSessionMsg();
try { // try {
FromDeviceMsg deviceMsg = sessionMsg.getMsg(); // FromDeviceMsg deviceMsg = sessionMsg.getMsg();
ToDeviceMsg toDeviceMsg = null; // ToDeviceMsg toDeviceMsg = null;
if (deviceMsg.getMsgType() == SessionMsgType.POST_TELEMETRY_REQUEST) { // if (deviceMsg.getMsgType() == SessionMsgType.POST_TELEMETRY_REQUEST) {
toDeviceMsg = BasicStatusCodeResponse.onSuccess(deviceMsg.getMsgType(), BasicRequest.DEFAULT_REQUEST_ID); // toDeviceMsg = BasicStatusCodeResponse.onSuccess(deviceMsg.getMsgType(), BasicRequest.DEFAULT_REQUEST_ID);
} else if (deviceMsg.getMsgType() == SessionMsgType.GET_ATTRIBUTES_REQUEST) { // } else if (deviceMsg.getMsgType() == SessionMsgType.GET_ATTRIBUTES_REQUEST) {
List<AttributeKvEntry> data = new ArrayList<>(); // List<AttributeKvEntry> data = new ArrayList<>();
data.add(new BaseAttributeKvEntry(new StringDataEntry("key1", "value1"), System.currentTimeMillis())); // data.add(new BaseAttributeKvEntry(new StringDataEntry("key1", "value1"), System.currentTimeMillis()));
data.add(new BaseAttributeKvEntry(new LongDataEntry("key2", 42L), System.currentTimeMillis())); // data.add(new BaseAttributeKvEntry(new LongDataEntry("key2", 42L), System.currentTimeMillis()));
BasicAttributeKVMsg kv = BasicAttributeKVMsg.fromClient(data); // BasicAttributeKVMsg kv = BasicAttributeKVMsg.fromClient(data);
toDeviceMsg = BasicGetAttributesResponse.onSuccess(deviceMsg.getMsgType(), BasicRequest.DEFAULT_REQUEST_ID, kv); // toDeviceMsg = BasicGetAttributesResponse.onSuccess(deviceMsg.getMsgType(), BasicRequest.DEFAULT_REQUEST_ID, kv);
} // }
if (toDeviceMsg != null) { // if (toDeviceMsg != null) {
sessionMsg.getSessionContext().onMsg(new BasicSessionActorToAdaptorMsg(sessionMsg.getSessionContext(), toDeviceMsg)); // sessionMsg.getSessionContext().onMsg(new BasicSessionActorToAdaptorMsg(sessionMsg.getSessionContext(), toDeviceMsg));
} // }
} catch (Exception e) { // } catch (Exception e) {
e.printStackTrace(); // e.printStackTrace();
} // }
} // }
} // }
@Override @Override
public void onDeviceAdded(Device device) { public void onDeviceAdded(Device device) {

302
transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -26,18 +26,25 @@ import io.netty.handler.codec.mqtt.MqttFixedHeader;
import io.netty.handler.codec.mqtt.MqttMessage; import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader; import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader;
import io.netty.handler.codec.mqtt.MqttPubAckMessage; import io.netty.handler.codec.mqtt.MqttPubAckMessage;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttQoS;
import io.netty.handler.codec.mqtt.MqttSubAckMessage; import io.netty.handler.codec.mqtt.MqttSubAckMessage;
import io.netty.handler.codec.mqtt.MqttSubAckPayload; import io.netty.handler.codec.mqtt.MqttSubAckPayload;
import io.netty.handler.codec.mqtt.MqttSubscribeMessage;
import io.netty.handler.codec.mqtt.MqttTopicSubscription;
import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage;
import io.netty.handler.ssl.SslHandler; import io.netty.handler.ssl.SslHandler;
import io.netty.util.concurrent.Future; import io.netty.util.concurrent.Future;
import io.netty.util.concurrent.GenericFutureListener; import io.netty.util.concurrent.GenericFutureListener;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.quota.QuotaService; import org.thingsboard.server.common.transport.quota.QuotaService;
import org.thingsboard.server.dao.EncryptionUtil; import org.thingsboard.server.dao.EncryptionUtil;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent; import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEventMsg; import org.thingsboard.server.gen.transport.TransportProtos.SessionEventMsg;
@ -54,6 +61,7 @@ import javax.net.ssl.SSLPeerUnverifiedException;
import javax.security.cert.X509Certificate; import javax.security.cert.X509Certificate;
import java.io.IOException; import java.io.IOException;
import java.net.InetSocketAddress; import java.net.InetSocketAddress;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -65,14 +73,16 @@ import static io.netty.handler.codec.mqtt.MqttConnectReturnCode.CONNECTION_REFUS
import static io.netty.handler.codec.mqtt.MqttMessageType.CONNACK; import static io.netty.handler.codec.mqtt.MqttMessageType.CONNACK;
import static io.netty.handler.codec.mqtt.MqttMessageType.PUBACK; import static io.netty.handler.codec.mqtt.MqttMessageType.PUBACK;
import static io.netty.handler.codec.mqtt.MqttMessageType.SUBACK; import static io.netty.handler.codec.mqtt.MqttMessageType.SUBACK;
import static io.netty.handler.codec.mqtt.MqttMessageType.UNSUBACK;
import static io.netty.handler.codec.mqtt.MqttQoS.AT_LEAST_ONCE; import static io.netty.handler.codec.mqtt.MqttQoS.AT_LEAST_ONCE;
import static io.netty.handler.codec.mqtt.MqttQoS.AT_MOST_ONCE; import static io.netty.handler.codec.mqtt.MqttQoS.AT_MOST_ONCE;
import static io.netty.handler.codec.mqtt.MqttQoS.FAILURE;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
@Slf4j @Slf4j
public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener<Future<? super Void>> { public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener<Future<? super Void>>, SessionMsgListener {
public static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; public static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE;
@ -84,8 +94,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final SslHandler sslHandler; private final SslHandler sslHandler;
private final ConcurrentMap<String, Integer> mqttQoSMap; private final ConcurrentMap<String, Integer> mqttQoSMap;
private final SessionInfoProto sessionInfo; private volatile SessionInfoProto sessionInfo;
private volatile InetSocketAddress address; private volatile InetSocketAddress address;
private volatile DeviceSessionCtx deviceSessionCtx; private volatile DeviceSessionCtx deviceSessionCtx;
private volatile GatewaySessionCtx gatewaySessionCtx; private volatile GatewaySessionCtx gatewaySessionCtx;
@ -98,11 +107,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
this.quotaService = context.getQuotaService(); this.quotaService = context.getQuotaService();
this.sslHandler = context.getSslHandler(); this.sslHandler = context.getSslHandler();
this.mqttQoSMap = new ConcurrentHashMap<>(); this.mqttQoSMap = new ConcurrentHashMap<>();
this.sessionInfo = SessionInfoProto.newBuilder()
.setNodeId(context.getNodeId())
.setSessionIdMSB(sessionId.getMostSignificantBits())
.setSessionIdLSB(sessionId.getLeastSignificantBits())
.build();
this.deviceSessionCtx = new DeviceSessionCtx(mqttQoSMap); this.deviceSessionCtx = new DeviceSessionCtx(mqttQoSMap);
} }
@ -135,15 +139,15 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case CONNECT: case CONNECT:
processConnect(ctx, (MqttConnectMessage) msg); processConnect(ctx, (MqttConnectMessage) msg);
break; break;
// case PUBLISH: case PUBLISH:
// processPublish(ctx, (MqttPublishMessage) msg); processPublish(ctx, (MqttPublishMessage) msg);
// break; break;
// case SUBSCRIBE: case SUBSCRIBE:
// processSubscribe(ctx, (MqttSubscribeMessage) msg); processSubscribe(ctx, (MqttSubscribeMessage) msg);
// break; break;
// case UNSUBSCRIBE: case UNSUBSCRIBE:
// processUnsubscribe(ctx, (MqttUnsubscribeMessage) msg); processUnsubscribe(ctx, (MqttUnsubscribeMessage) msg);
// break; break;
// case PINGREQ: // case PINGREQ:
// if (checkConnected(ctx)) { // if (checkConnected(ctx)) {
// ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); // ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0)));
@ -160,24 +164,25 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
// private void processPublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg) { private void processPublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg) {
// if (!checkConnected(ctx)) { if (!checkConnected(ctx)) {
// return; return;
// } }
// String topicName = mqttMsg.variableHeader().topicName(); String topicName = mqttMsg.variableHeader().topicName();
// int msgId = mqttMsg.variableHeader().packetId(); int msgId = mqttMsg.variableHeader().packetId();
// log.trace("[{}] Processing publish msg [{}][{}]!", sessionId, topicName, msgId); log.trace("[{}] Processing publish msg [{}][{}]!", sessionId, topicName, msgId);
//
// if (topicName.startsWith(BASE_GATEWAY_API_TOPIC)) { if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) {
// if (gatewaySessionCtx != null) { if (gatewaySessionCtx != null) {
// gatewaySessionCtx.setChannel(ctx); gatewaySessionCtx.setChannel(ctx);
// handleMqttPublishMsg(topicName, msgId, mqttMsg); // handleMqttPublishMsg(topicName, msgId, mqttMsg);
// } }
// } else { } else {
// processDevicePublish(ctx, mqttMsg, topicName, msgId); processDevicePublish(ctx, mqttMsg, topicName, msgId);
// } }
// } }
//
//
// private void handleMqttPublishMsg(String topicName, int msgId, MqttPublishMessage mqttMsg) { // private void handleMqttPublishMsg(String topicName, int msgId, MqttPublishMessage mqttMsg) {
// try { // try {
// switch (topicName) { // switch (topicName) {
@ -205,7 +210,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
// } // }
// } // }
// //
// private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) { private void processDevicePublish(ChannelHandlerContext ctx, MqttPublishMessage mqttMsg, String topicName, int msgId) {
try {
if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_TOPIC)) {
TransportProtos.PostTelemetryMsg postTelemetryMsg = adaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg);
transportService.process(sessionInfo, postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg));
} else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_TOPIC)) {
TransportProtos.PostAttributeMsg postAttributeMsg = adaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg);
transportService.process(sessionInfo, postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg));
} else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) {
TransportProtos.GetAttributeRequestMsg getAttributeMsg = adaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg);
transportService.process(sessionInfo, getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg));
}
} catch (AdaptorException e) {
log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e);
log.info("[{}] Closing current session due to invalid publish msg [{}][{}]", sessionId, topicName, msgId);
ctx.close();
}
// AdaptorToSessionActorMsg msg = null; // AdaptorToSessionActorMsg msg = null;
// try { // try {
// if (topicName.equals(DEVICE_TELEMETRY_TOPIC)) { // if (topicName.equals(DEVICE_TELEMETRY_TOPIC)) {
@ -237,20 +258,38 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
// log.info("[{}] Closing current session due to invalid publish msg [{}][{}]", sessionId, topicName, msgId); // log.info("[{}] Closing current session due to invalid publish msg [{}][{}]", sessionId, topicName, msgId);
// ctx.close(); // ctx.close();
// } // }
// } }
//
// private void processSubscribe(ChannelHandlerContext ctx, MqttSubscribeMessage mqttMsg) { private <T> TransportServiceCallback<Void> getPubAckCallback(final ChannelHandlerContext ctx, final int msgId, final T msg) {
// if (!checkConnected(ctx)) { return new TransportServiceCallback<Void>() {
// return; @Override
// } public void onSuccess(Void dummy) {
// log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); log.trace("[{}] Published msg: {}", sessionId, msg);
// List<Integer> grantedQoSList = new ArrayList<>(); if (msgId > 0) {
// for (MqttTopicSubscription subscription : mqttMsg.payload().topicSubscriptions()) { ctx.writeAndFlush(createMqttPubAckMsg(msgId));
// String topic = subscription.topicName(); }
// MqttQoS reqQoS = subscription.qualityOfService(); }
// try {
// switch (topic) { @Override
// case DEVICE_ATTRIBUTES_TOPIC: { public void onError(Throwable e) {
log.trace("[{}] Failed to publish msg: {}", sessionId, msg, e);
ctx.close();
}
};
}
private void processSubscribe(ChannelHandlerContext ctx, MqttSubscribeMessage mqttMsg) {
if (!checkConnected(ctx)) {
return;
}
log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId());
List<Integer> grantedQoSList = new ArrayList<>();
for (MqttTopicSubscription subscription : mqttMsg.payload().topicSubscriptions()) {
String topic = subscription.topicName();
MqttQoS reqQoS = subscription.qualityOfService();
try {
switch (topic) {
// case MqttTopics.DEVICE_ATTRIBUTES_TOPIC: {
// AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, SUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); // AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, SUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg);
// processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg)); // processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg));
// registerSubQoS(topic, grantedQoSList, reqQoS); // registerSubQoS(topic, grantedQoSList, reqQoS);
@ -267,37 +306,37 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
// case GATEWAY_RPC_TOPIC: // case GATEWAY_RPC_TOPIC:
// registerSubQoS(topic, grantedQoSList, reqQoS); // registerSubQoS(topic, grantedQoSList, reqQoS);
// break; // break;
// case DEVICE_ATTRIBUTES_RESPONSES_TOPIC: case MqttTopics.DEVICE_ATTRIBUTES_RESPONSES_TOPIC:
// deviceSessionCtx.setAllowAttributeResponses(); deviceSessionCtx.setAllowAttributeResponses();
// registerSubQoS(topic, grantedQoSList, reqQoS); registerSubQoS(topic, grantedQoSList, reqQoS);
// break; break;
// default: default:
// log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topic, reqQoS); log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topic, reqQoS);
// grantedQoSList.add(FAILURE.value()); grantedQoSList.add(FAILURE.value());
// break; break;
// } }
// } catch (AdaptorException e) { } catch (Exception e) {
// log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topic, reqQoS); log.warn("[{}] Failed to subscribe to [{}][{}]", sessionId, topic, reqQoS);
// grantedQoSList.add(FAILURE.value()); grantedQoSList.add(FAILURE.value());
// } }
// } }
// ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList));
// } }
//
// private void registerSubQoS(String topic, List<Integer> grantedQoSList, MqttQoS reqQoS) { private void registerSubQoS(String topic, List<Integer> grantedQoSList, MqttQoS reqQoS) {
// grantedQoSList.add(getMinSupportedQos(reqQoS)); grantedQoSList.add(getMinSupportedQos(reqQoS));
// mqttQoSMap.put(topic, getMinSupportedQos(reqQoS)); mqttQoSMap.put(topic, getMinSupportedQos(reqQoS));
// } }
//
// private void processUnsubscribe(ChannelHandlerContext ctx, MqttUnsubscribeMessage mqttMsg) { private void processUnsubscribe(ChannelHandlerContext ctx, MqttUnsubscribeMessage mqttMsg) {
// if (!checkConnected(ctx)) { if (!checkConnected(ctx)) {
// return; return;
// } }
// log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId()); log.trace("[{}] Processing subscription [{}]!", sessionId, mqttMsg.variableHeader().messageId());
// for (String topicName : mqttMsg.payload().topics()) { for (String topicName : mqttMsg.payload().topics()) {
// mqttQoSMap.remove(topicName); mqttQoSMap.remove(topicName);
// try { try {
// switch (topicName) { switch (topicName) {
// case DEVICE_ATTRIBUTES_TOPIC: { // case DEVICE_ATTRIBUTES_TOPIC: {
// AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, UNSUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg); // AdaptorToSessionActorMsg msg = adaptor.convertToActorMsg(deviceSessionCtx, UNSUBSCRIBE_ATTRIBUTES_REQUEST, mqttMsg);
// processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg)); // processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg));
@ -308,23 +347,23 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
// processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg)); // processor.process(new BasicTransportToDeviceSessionActorMsg(deviceSessionCtx.getDevice(), msg));
// break; // break;
// } // }
// case DEVICE_ATTRIBUTES_RESPONSES_TOPIC: case MqttTopics.DEVICE_ATTRIBUTES_RESPONSES_TOPIC:
// deviceSessionCtx.setDisallowAttributeResponses(); deviceSessionCtx.setDisallowAttributeResponses();
// break; break;
// } }
// } catch (AdaptorException e) { } catch (Exception e) {
// log.warn("[{}] Failed to process unsubscription [{}] to [{}]", sessionId, mqttMsg.variableHeader().messageId(), topicName); log.warn("[{}] Failed to process unsubscription [{}] to [{}]", sessionId, mqttMsg.variableHeader().messageId(), topicName);
// } }
// } }
// ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId())); ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId()));
// } }
//
// private MqttMessage createUnSubAckMessage(int msgId) { private MqttMessage createUnSubAckMessage(int msgId) {
// MqttFixedHeader mqttFixedHeader = MqttFixedHeader mqttFixedHeader =
// new MqttFixedHeader(UNSUBACK, false, AT_LEAST_ONCE, false, 0); new MqttFixedHeader(UNSUBACK, false, AT_LEAST_ONCE, false, 0);
// MqttMessageIdVariableHeader mqttMessageIdVariableHeader = MqttMessageIdVariableHeader.from(msgId); MqttMessageIdVariableHeader mqttMessageIdVariableHeader = MqttMessageIdVariableHeader.from(msgId);
// return new MqttMessage(mqttFixedHeader, mqttMessageIdVariableHeader); return new MqttMessage(mqttFixedHeader, mqttMessageIdVariableHeader);
// } }
private void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { private void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) {
log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier()); log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier());
@ -346,15 +385,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
new TransportServiceCallback<ValidateDeviceCredentialsResponseMsg>() { new TransportServiceCallback<ValidateDeviceCredentialsResponseMsg>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) {
if (!msg.hasDeviceInfo()) { onValidateDeviceResponse(msg, ctx);
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED));
ctx.close();
} else {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED));
deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo());
transportService.process(deviceSessionCtx, getSessionEventMsg(SessionEvent.OPEN), null);
checkGatewaySession();
}
} }
@Override @Override
@ -375,15 +406,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
new TransportServiceCallback<ValidateDeviceCredentialsResponseMsg>() { new TransportServiceCallback<ValidateDeviceCredentialsResponseMsg>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) {
if (!msg.hasDeviceInfo()) { onValidateDeviceResponse(msg, ctx);
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED));
ctx.close();
} else {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED));
deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo());
transportService.process(deviceSessionCtx, getSessionEventMsg(SessionEvent.OPEN), null);
checkGatewaySession();
}
} }
@Override @Override
@ -415,7 +438,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private void processDisconnect(ChannelHandlerContext ctx) { private void processDisconnect(ChannelHandlerContext ctx) {
ctx.close(); ctx.close();
if (deviceSessionCtx.isConnected()) { if (deviceSessionCtx.isConnected()) {
transportService.process(deviceSessionCtx, getSessionEventMsg(SessionEvent.CLOSED), null); transportService.process(sessionInfo, getSessionEventMsg(SessionEvent.CLOSED), null);
transportService.deregisterSession(sessionInfo);
if (gatewaySessionCtx != null) { if (gatewaySessionCtx != null) {
gatewaySessionCtx.onGatewayDisconnect(); gatewaySessionCtx.onGatewayDisconnect();
} }
@ -488,16 +512,46 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private SessionEventMsg getSessionEventMsg(SessionEvent event) { private SessionEventMsg getSessionEventMsg(SessionEvent event) {
return SessionEventMsg.newBuilder() return SessionEventMsg.newBuilder()
.setSessionInfo(sessionInfo) .setSessionType(TransportProtos.SessionType.ASYNC)
.setDeviceIdMSB(deviceSessionCtx.getDeviceIdMSB())
.setDeviceIdLSB(deviceSessionCtx.getDeviceIdLSB())
.setEvent(event).build(); .setEvent(event).build();
} }
@Override @Override
public void operationComplete(Future<? super Void> future) throws Exception { public void operationComplete(Future<? super Void> future) throws Exception {
if (deviceSessionCtx.isConnected()) { if (deviceSessionCtx.isConnected()) {
transportService.process(deviceSessionCtx, getSessionEventMsg(SessionEvent.CLOSED), null); transportService.process(sessionInfo, getSessionEventMsg(SessionEvent.CLOSED), null);
transportService.deregisterSession(sessionInfo);
}
}
private void onValidateDeviceResponse(ValidateDeviceCredentialsResponseMsg msg, ChannelHandlerContext ctx) {
if (!msg.hasDeviceInfo()) {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED));
ctx.close();
} else {
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED));
deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo());
sessionInfo = SessionInfoProto.newBuilder()
.setNodeId(context.getNodeId())
.setSessionIdMSB(sessionId.getMostSignificantBits())
.setSessionIdLSB(sessionId.getLeastSignificantBits())
.setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB())
.setDeviceIdLSB(msg.getDeviceInfo().getDeviceIdLSB())
.setTenantIdMSB(msg.getDeviceInfo().getTenantIdMSB())
.setTenantIdLSB(msg.getDeviceInfo().getTenantIdLSB())
.build();
transportService.process(sessionInfo, getSessionEventMsg(SessionEvent.OPEN), null);
transportService.registerSession(sessionInfo, this);
checkGatewaySession();
}
}
@Override
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg response) {
try {
adaptor.convertToPublish(deviceSessionCtx, response).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush);
} catch (Exception e) {
log.trace("[{}] Failed to convert device attributes to MQTT msg", sessionId, e);
} }
} }
} }

61
transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java

@ -17,6 +17,7 @@ package org.thingsboard.server.transport.mqtt.adaptors;
import com.google.gson.Gson; import com.google.gson.Gson;
import com.google.gson.JsonElement; import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser; import com.google.gson.JsonParser;
import com.google.gson.JsonSyntaxException; import com.google.gson.JsonSyntaxException;
import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBuf;
@ -25,12 +26,14 @@ import io.netty.buffer.UnpooledByteBufAllocator;
import io.netty.handler.codec.mqtt.*; import io.netty.handler.codec.mqtt.*;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.id.SessionId; import org.thingsboard.server.common.data.id.SessionId;
import org.thingsboard.server.common.msg.core.*; import org.thingsboard.server.common.msg.core.*;
import org.thingsboard.server.common.msg.kv.AttributesKVMsg; import org.thingsboard.server.common.msg.kv.AttributesKVMsg;
import org.thingsboard.server.common.msg.session.*; import org.thingsboard.server.common.msg.session.*;
import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.mqtt.MqttTopics; import org.thingsboard.server.transport.mqtt.MqttTopics;
import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx;
import org.thingsboard.server.transport.mqtt.MqttTransportHandler; import org.thingsboard.server.transport.mqtt.MqttTransportHandler;
@ -52,6 +55,64 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
private static final Charset UTF8 = Charset.forName("UTF-8"); private static final Charset UTF8 = Charset.forName("UTF-8");
private static final ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); private static final ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false);
@Override
public TransportProtos.PostTelemetryMsg convertToPostTelemetry(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException {
String payload = validatePayload(ctx.getSessionId(), inbound.payload());
try {
return JsonConverter.convertToTelemetryProto(new JsonParser().parse(payload));
} catch (IllegalStateException | JsonSyntaxException ex) {
throw new AdaptorException(ex);
}
}
@Override
public TransportProtos.PostAttributeMsg convertToPostAttributes(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException {
String payload = validatePayload(ctx.getSessionId(), inbound.payload());
try {
return JsonConverter.convertToAttributesProto(new JsonParser().parse(payload));
} catch (IllegalStateException | JsonSyntaxException ex) {
throw new AdaptorException(ex);
}
}
@Override
public TransportProtos.GetAttributeRequestMsg convertToGetAttributes(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException {
String topicName = inbound.variableHeader().topicName();
try {
TransportProtos.GetAttributeRequestMsg.Builder result = TransportProtos.GetAttributeRequestMsg.newBuilder();
result.setRequestId(Integer.valueOf(topicName.substring(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX.length())));
String payload = inbound.payload().toString(UTF8);
JsonElement requestBody = new JsonParser().parse(payload);
Set<String> clientKeys = toStringSet(requestBody, "clientKeys");
Set<String> sharedKeys = toStringSet(requestBody, "sharedKeys");
if (clientKeys != null) {
result.addAllClientAttributeNames(clientKeys);
}
if (sharedKeys != null) {
result.addAllSharedAttributeNames(sharedKeys);
}
return result.build();
} catch (RuntimeException e) {
log.warn("Failed to decode get attributes request", e);
throw new AdaptorException(e);
}
}
@Override
public Optional<MqttMessage> convertToPublish(DeviceSessionCtx ctx, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException {
if (!StringUtils.isEmpty(responseMsg.getError())) {
throw new AdaptorException(responseMsg.getError());
} else {
Integer requestId = responseMsg.getRequestId();
if (requestId >= 0) {
return Optional.of(createMqttPublishMsg(ctx,
MqttTopics.DEVICE_ATTRIBUTES_RESPONSE_TOPIC_PREFIX + requestId,
JsonConverter.toJson(responseMsg)));
}
return Optional.empty();
}
}
@Override @Override
public AdaptorToSessionActorMsg convertToActorMsg(DeviceSessionCtx ctx, SessionMsgType type, MqttMessage inbound) throws AdaptorException { public AdaptorToSessionActorMsg convertToActorMsg(DeviceSessionCtx ctx, SessionMsgType type, MqttMessage inbound) throws AdaptorException {
FromDeviceMsg msg; FromDeviceMsg msg;

13
transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java

@ -16,11 +16,24 @@
package org.thingsboard.server.transport.mqtt.adaptors; package org.thingsboard.server.transport.mqtt.adaptors;
import io.netty.handler.codec.mqtt.MqttMessage; import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import org.thingsboard.server.common.transport.TransportAdaptor; import org.thingsboard.server.common.transport.TransportAdaptor;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx; import org.thingsboard.server.transport.mqtt.session.DeviceSessionCtx;
import java.util.Optional;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
public interface MqttTransportAdaptor extends TransportAdaptor<DeviceSessionCtx, MqttMessage, MqttMessage> { public interface MqttTransportAdaptor extends TransportAdaptor<DeviceSessionCtx, MqttMessage, MqttMessage> {
TransportProtos.PostTelemetryMsg convertToPostTelemetry(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException;
TransportProtos.PostAttributeMsg convertToPostAttributes(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException;
TransportProtos.GetAttributeRequestMsg convertToGetAttributes(DeviceSessionCtx ctx, MqttPublishMessage inbound) throws AdaptorException;
Optional<MqttMessage> convertToPublish(DeviceSessionCtx ctx, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException;
} }

4
transport/mqtt-common/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java

@ -17,6 +17,7 @@ package org.thingsboard.server.transport.mqtt.session;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.mqtt.*; import io.netty.handler.codec.mqtt.*;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.id.SessionId; import org.thingsboard.server.common.data.id.SessionId;
import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg; import org.thingsboard.server.common.msg.session.SessionActorToAdaptorMsg;
@ -41,12 +42,13 @@ import java.util.concurrent.atomic.AtomicInteger;
public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
private final MqttSessionId sessionId; private final MqttSessionId sessionId;
@Getter
private ChannelHandlerContext channel; private ChannelHandlerContext channel;
private volatile boolean allowAttributeResponses; private volatile boolean allowAttributeResponses;
private AtomicInteger msgIdSeq = new AtomicInteger(0); private AtomicInteger msgIdSeq = new AtomicInteger(0);
public DeviceSessionCtx(ConcurrentMap<String, Integer> mqttQoSMap) { public DeviceSessionCtx(ConcurrentMap<String, Integer> mqttQoSMap) {
super(null, null, null); super(null, null, mqttQoSMap);
this.sessionId = new MqttSessionId(); this.sessionId = new MqttSessionId();
} }

33
transport/mqtt-transport/src/main/java/org/thingsboard/server/mqtt/service/MqttTransportService.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@ -111,8 +111,8 @@ public class MqttTransportService implements TransportService {
TBKafkaConsumerTemplate.TBKafkaConsumerTemplateBuilder<TransportApiResponseMsg> responseBuilder = TBKafkaConsumerTemplate.builder(); TBKafkaConsumerTemplate.TBKafkaConsumerTemplateBuilder<TransportApiResponseMsg> responseBuilder = TBKafkaConsumerTemplate.builder();
responseBuilder.settings(kafkaSettings); responseBuilder.settings(kafkaSettings);
responseBuilder.topic(transportApiResponsesTopic + "." + transportContext.getNodeId()); responseBuilder.topic(transportApiResponsesTopic + "." + transportContext.getNodeId());
responseBuilder.clientId(transportContext.getNodeId()); responseBuilder.clientId("transport-api-client-" + transportContext.getNodeId());
responseBuilder.groupId(null); responseBuilder.groupId("transport-api-client");
responseBuilder.autoCommit(true); responseBuilder.autoCommit(true);
responseBuilder.autoCommitIntervalMs(autoCommitInterval); responseBuilder.autoCommitIntervalMs(autoCommitInterval);
responseBuilder.decoder(new TransportApiResponseDecoder()); responseBuilder.decoder(new TransportApiResponseDecoder());
@ -137,8 +137,8 @@ public class MqttTransportService implements TransportService {
TBKafkaConsumerTemplate.TBKafkaConsumerTemplateBuilder<ToTransportMsg> mainConsumerBuilder = TBKafkaConsumerTemplate.builder(); TBKafkaConsumerTemplate.TBKafkaConsumerTemplateBuilder<ToTransportMsg> mainConsumerBuilder = TBKafkaConsumerTemplate.builder();
mainConsumerBuilder.settings(kafkaSettings); mainConsumerBuilder.settings(kafkaSettings);
mainConsumerBuilder.topic(notificationsTopic + "." + transportContext.getNodeId()); mainConsumerBuilder.topic(notificationsTopic + "." + transportContext.getNodeId());
mainConsumerBuilder.clientId(transportContext.getNodeId()); mainConsumerBuilder.clientId("transport-" + transportContext.getNodeId());
mainConsumerBuilder.groupId(null); mainConsumerBuilder.groupId("transport");
mainConsumerBuilder.autoCommit(true); mainConsumerBuilder.autoCommit(true);
mainConsumerBuilder.autoCommitIntervalMs(notificationsAutoCommitInterval); mainConsumerBuilder.autoCommitIntervalMs(notificationsAutoCommitInterval);
mainConsumerBuilder.decoder(new ToTransportMsgResponseDecoder()); mainConsumerBuilder.decoder(new ToTransportMsgResponseDecoder());
@ -242,6 +242,15 @@ public class MqttTransportService implements TransportService {
send(sessionInfo, toRuleEngineMsg, callback); send(sessionInfo, toRuleEngineMsg, callback);
} }
@Override
public void process(SessionInfoProto sessionInfo, TransportProtos.GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback) {
ToRuleEngineMsg toRuleEngineMsg = ToRuleEngineMsg.newBuilder().setToDeviceActorMsg(
TransportProtos.TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setGetAttributes(msg).build()
).build();
send(sessionInfo, toRuleEngineMsg, callback);
}
@Override @Override
public void registerSession(SessionInfoProto sessionInfo, SessionMsgListener listener) { public void registerSession(SessionInfoProto sessionInfo, SessionMsgListener listener) {
sessions.putIfAbsent(toId(sessionInfo), listener); sessions.putIfAbsent(toId(sessionInfo), listener);
@ -271,9 +280,13 @@ public class MqttTransportService implements TransportService {
@Override @Override
public void onCompletion(RecordMetadata metadata, Exception exception) { public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception == null) { if (exception == null) {
callback.onSuccess(null); if (callback != null) {
callback.onSuccess(null);
}
} else { } else {
callback.onError(exception); if (callback != null) {
callback.onError(exception);
}
} }
} }
} }

8
transport/mqtt-transport/src/main/java/org/thingsboard/server/mqtt/service/ToRuleEngineMsgEncoder.java

@ -1,12 +1,12 @@
/** /**
* Copyright © 2016-2018 The Thingsboard Authors * Copyright © 2016-2018 The Thingsboard Authors
* <p> *
* Licensed under the Apache License, Version 2.0 (the "License"); * Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* <p> *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* <p> *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.

Loading…
Cancel
Save