|
|
|
@ -18,31 +18,63 @@ package org.thingsboard.server.actors.device; |
|
|
|
import akka.actor.ActorContext; |
|
|
|
import akka.actor.ActorRef; |
|
|
|
import akka.event.LoggingAdapter; |
|
|
|
import com.datastax.driver.core.utils.UUIDs; |
|
|
|
import com.google.gson.Gson; |
|
|
|
import com.google.gson.JsonArray; |
|
|
|
import com.google.gson.JsonObject; |
|
|
|
import org.thingsboard.server.actors.ActorSystemContext; |
|
|
|
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
import org.thingsboard.server.common.data.Device; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.SessionId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKey; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.KvEntry; |
|
|
|
import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; |
|
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
|
import org.thingsboard.server.common.msg.TbMsgDataType; |
|
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|
|
|
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|
|
|
import org.thingsboard.server.common.msg.core.*; |
|
|
|
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
|
|
import org.thingsboard.server.common.msg.core.AttributesUpdateNotification; |
|
|
|
import org.thingsboard.server.common.msg.core.BasicCommandAckResponse; |
|
|
|
import org.thingsboard.server.common.msg.core.BasicStatusCodeResponse; |
|
|
|
import org.thingsboard.server.common.msg.core.BasicToDeviceSessionActorMsg; |
|
|
|
import org.thingsboard.server.common.msg.core.RuleEngineError; |
|
|
|
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.ToDeviceRpcResponseMsg; |
|
|
|
import org.thingsboard.server.common.msg.core.ToDeviceSessionActorMsg; |
|
|
|
import org.thingsboard.server.common.msg.device.DeviceToDeviceActorMsg; |
|
|
|
import org.thingsboard.server.common.msg.kv.BasicAttributeKVMsg; |
|
|
|
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; |
|
|
|
import org.thingsboard.server.common.msg.session.FromDeviceMsg; |
|
|
|
import org.thingsboard.server.common.msg.session.MsgType; |
|
|
|
import org.thingsboard.server.common.msg.session.FromDeviceRequestMsg; |
|
|
|
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.ToDeviceMsg; |
|
|
|
import org.thingsboard.server.extensions.api.device.DeviceAttributes; |
|
|
|
import org.thingsboard.server.common.msg.timeout.DeviceActorQueueTimeoutMsg; |
|
|
|
import org.thingsboard.server.common.msg.timeout.DeviceActorRpcTimeoutMsg; |
|
|
|
import org.thingsboard.server.extensions.api.device.DeviceAttributesEventNotificationMsg; |
|
|
|
import org.thingsboard.server.extensions.api.device.DeviceNameOrTypeUpdateMsg; |
|
|
|
import org.thingsboard.server.extensions.api.plugins.msg.*; |
|
|
|
|
|
|
|
import java.util.*; |
|
|
|
import org.thingsboard.server.extensions.api.plugins.msg.FromDeviceRpcResponse; |
|
|
|
import org.thingsboard.server.extensions.api.plugins.msg.RpcError; |
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.HashMap; |
|
|
|
import java.util.HashSet; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ExecutionException; |
|
|
|
import java.util.concurrent.TimeoutException; |
|
|
|
import java.util.function.Consumer; |
|
|
|
@ -54,25 +86,30 @@ import java.util.stream.Collectors; |
|
|
|
*/ |
|
|
|
public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { |
|
|
|
|
|
|
|
private final TenantId tenantId; |
|
|
|
private final DeviceId deviceId; |
|
|
|
private final Map<SessionId, SessionInfo> sessions; |
|
|
|
private final Map<SessionId, SessionInfo> attributeSubscriptions; |
|
|
|
private final Map<SessionId, SessionInfo> rpcSubscriptions; |
|
|
|
|
|
|
|
private final Map<Integer, ToDeviceRpcRequestMetadata> rpcPendingMap; |
|
|
|
private final Map<UUID, PendingSessionMsgData> pendingMsgs; |
|
|
|
|
|
|
|
private final Gson gson = new Gson(); |
|
|
|
|
|
|
|
private int rpcSeq = 0; |
|
|
|
private String deviceName; |
|
|
|
private String deviceType; |
|
|
|
private DeviceAttributes deviceAttributes; |
|
|
|
private TbMsgMetaData defaultMetaData; |
|
|
|
|
|
|
|
public DeviceActorMessageProcessor(ActorSystemContext systemContext, LoggingAdapter logger, DeviceId deviceId) { |
|
|
|
public DeviceActorMessageProcessor(ActorSystemContext systemContext, LoggingAdapter logger, TenantId tenantId, DeviceId deviceId) { |
|
|
|
super(systemContext, logger); |
|
|
|
this.tenantId = tenantId; |
|
|
|
this.deviceId = deviceId; |
|
|
|
this.sessions = new HashMap<>(); |
|
|
|
this.attributeSubscriptions = new HashMap<>(); |
|
|
|
this.rpcSubscriptions = new HashMap<>(); |
|
|
|
this.rpcPendingMap = new HashMap<>(); |
|
|
|
this.pendingMsgs = new HashMap<>(); |
|
|
|
initAttributes(); |
|
|
|
} |
|
|
|
|
|
|
|
@ -81,19 +118,12 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
Device device = systemContext.getDeviceService().findDeviceById(deviceId); |
|
|
|
this.deviceName = device.getName(); |
|
|
|
this.deviceType = device.getType(); |
|
|
|
this.deviceAttributes = new DeviceAttributes(fetchAttributes(DataConstants.CLIENT_SCOPE), |
|
|
|
fetchAttributes(DataConstants.SERVER_SCOPE), fetchAttributes(DataConstants.SHARED_SCOPE)); |
|
|
|
this.defaultMetaData = new TbMsgMetaData(); |
|
|
|
this.defaultMetaData.putValue("deviceName", deviceName); |
|
|
|
this.defaultMetaData.putValue("deviceType", deviceType); |
|
|
|
} |
|
|
|
|
|
|
|
private void refreshAttributes(DeviceAttributesEventNotificationMsg msg) { |
|
|
|
if (msg.isDeleted()) { |
|
|
|
msg.getDeletedKeys().forEach(key -> deviceAttributes.remove(key)); |
|
|
|
} else { |
|
|
|
deviceAttributes.update(msg.getScope(), msg.getValues()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
void processRpcRequest(ActorContext context, ToDeviceRpcRequestPluginMsg msg) { |
|
|
|
void processRpcRequest(ActorContext context, org.thingsboard.server.service.rpc.ToDeviceRpcRequestMsg msg) { |
|
|
|
ToDeviceRpcRequest request = msg.getMsg(); |
|
|
|
ToDeviceRpcRequestBody body = request.getBody(); |
|
|
|
ToDeviceRpcRequestMsg rpcRequest = new ToDeviceRpcRequestMsg( |
|
|
|
@ -120,9 +150,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
syncSessionSet.forEach(rpcSubscriptions::remove); |
|
|
|
|
|
|
|
if (request.isOneway() && sent) { |
|
|
|
ToPluginRpcResponseDeviceMsg responsePluginMsg = toPluginRpcResponseMsg(msg, (String) null); |
|
|
|
context.parent().tell(responsePluginMsg, ActorRef.noSender()); |
|
|
|
logger.debug("[{}] Rpc command response sent [{}]!", deviceId, request.getId()); |
|
|
|
systemContext.getDeviceRpcService().process(new FromDeviceRpcResponse(msg.getMsg().getId(), null, null)); |
|
|
|
} else { |
|
|
|
registerPendingRpcRequest(context, msg, sent, rpcRequest, timeout); |
|
|
|
} |
|
|
|
@ -134,18 +163,36 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
private void registerPendingRpcRequest(ActorContext context, ToDeviceRpcRequestPluginMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { |
|
|
|
private void registerPendingRpcRequest(ActorContext context, org.thingsboard.server.service.rpc.ToDeviceRpcRequestMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { |
|
|
|
rpcPendingMap.put(rpcRequest.getRequestId(), new ToDeviceRpcRequestMetadata(msg, sent)); |
|
|
|
TimeoutIntMsg timeoutMsg = new TimeoutIntMsg(rpcRequest.getRequestId(), timeout); |
|
|
|
DeviceActorRpcTimeoutMsg timeoutMsg = new DeviceActorRpcTimeoutMsg(rpcRequest.getRequestId(), timeout); |
|
|
|
scheduleMsgWithDelay(context, timeoutMsg, timeoutMsg.getTimeout()); |
|
|
|
} |
|
|
|
|
|
|
|
public void processTimeout(ActorContext context, TimeoutMsg msg) { |
|
|
|
void processRpcTimeout(ActorContext context, DeviceActorRpcTimeoutMsg msg) { |
|
|
|
ToDeviceRpcRequestMetadata requestMd = rpcPendingMap.remove(msg.getId()); |
|
|
|
if (requestMd != null) { |
|
|
|
logger.debug("[{}] RPC request [{}] timeout detected!", deviceId, msg.getId()); |
|
|
|
ToPluginRpcResponseDeviceMsg responsePluginMsg = toPluginRpcResponseMsg(requestMd.getMsg(), requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION); |
|
|
|
context.parent().tell(responsePluginMsg, ActorRef.noSender()); |
|
|
|
systemContext.getDeviceRpcService().process(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), |
|
|
|
null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
void processQueueTimeout(ActorContext context, DeviceActorQueueTimeoutMsg msg) { |
|
|
|
PendingSessionMsgData data = pendingMsgs.remove(msg.getId()); |
|
|
|
if (data != null) { |
|
|
|
logger.debug("[{}] Queue put [{}] timeout detected!", deviceId, msg.getId()); |
|
|
|
ToDeviceMsg toDeviceMsg = new RuleEngineErrorMsg(data.getSessionMsgType(), RuleEngineError.QUEUE_PUT_TIMEOUT); |
|
|
|
sendMsgToSessionActor(new BasicToDeviceSessionActorMsg(toDeviceMsg, data.getSessionId()), data.getServerAddress()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
void processQueueAck(ActorContext context, RuleEngineQueuePutAckMsg msg) { |
|
|
|
PendingSessionMsgData data = pendingMsgs.remove(msg.getId()); |
|
|
|
if (data != null) { |
|
|
|
logger.debug("[{}] Queue put [{}] ack detected!", deviceId, msg.getId()); |
|
|
|
ToDeviceMsg toDeviceMsg = BasicStatusCodeResponse.onSuccess(data.getSessionMsgType(), data.getRequestId()); |
|
|
|
sendMsgToSessionActor(new BasicToDeviceSessionActorMsg(toDeviceMsg, data.getSessionId()), data.getServerAddress()); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -175,8 +222,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
ToDeviceRpcRequestBody body = request.getBody(); |
|
|
|
if (request.isOneway()) { |
|
|
|
sentOneWayIds.add(entry.getKey()); |
|
|
|
ToPluginRpcResponseDeviceMsg responsePluginMsg = toPluginRpcResponseMsg(entry.getValue().getMsg(), (String) null); |
|
|
|
context.parent().tell(responsePluginMsg, ActorRef.noSender()); |
|
|
|
systemContext.getDeviceRpcService().process(new FromDeviceRpcResponse(request.getId(), null, null)); |
|
|
|
} |
|
|
|
ToDeviceRpcRequestMsg rpcRequest = new ToDeviceRpcRequestMsg( |
|
|
|
entry.getKey(), |
|
|
|
@ -188,14 +234,70 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
void process(ActorContext context, ToDeviceActorMsg msg) { |
|
|
|
void process(ActorContext context, DeviceToDeviceActorMsg msg) { |
|
|
|
processSubscriptionCommands(context, msg); |
|
|
|
processRpcResponses(context, msg); |
|
|
|
processSessionStateMsgs(msg); |
|
|
|
SessionMsgType sessionMsgType = msg.getPayload().getMsgType(); |
|
|
|
if (sessionMsgType.requiresRulesProcessing()) { |
|
|
|
switch (sessionMsgType) { |
|
|
|
case GET_ATTRIBUTES_REQUEST: |
|
|
|
handleGetAttributesRequest(msg); |
|
|
|
break; |
|
|
|
case POST_ATTRIBUTES_REQUEST: |
|
|
|
break; |
|
|
|
case POST_TELEMETRY_REQUEST: |
|
|
|
handlePostTelemetryRequest(context, msg); |
|
|
|
break; |
|
|
|
case TO_SERVER_RPC_REQUEST: |
|
|
|
break; |
|
|
|
//TODO: push to queue and start processing!
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void handleGetAttributesRequest(DeviceToDeviceActorMsg msg) { |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
private void handlePostTelemetryRequest(ActorContext context, DeviceToDeviceActorMsg src) { |
|
|
|
TelemetryUploadRequest telemetry = (TelemetryUploadRequest) src.getPayload(); |
|
|
|
|
|
|
|
Map<Long, List<KvEntry>> tsData = telemetry.getData(); |
|
|
|
|
|
|
|
JsonArray json = new JsonArray(); |
|
|
|
for (Map.Entry<Long, List<KvEntry>> entry : tsData.entrySet()) { |
|
|
|
JsonObject ts = new JsonObject(); |
|
|
|
ts.addProperty("ts", entry.getKey()); |
|
|
|
JsonObject values = new JsonObject(); |
|
|
|
for (KvEntry kv : entry.getValue()) { |
|
|
|
kv.getBooleanValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); |
|
|
|
kv.getLongValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); |
|
|
|
kv.getDoubleValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); |
|
|
|
kv.getStrValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); |
|
|
|
} |
|
|
|
ts.add("values", values); |
|
|
|
json.add(ts); |
|
|
|
} |
|
|
|
|
|
|
|
TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, defaultMetaData, TbMsgDataType.JSON, gson.toJson(json)); |
|
|
|
pushToRuleEngineWithTimeout(context, tbMsg, src, telemetry); |
|
|
|
} |
|
|
|
|
|
|
|
private void pushToRuleEngineWithTimeout(ActorContext context, TbMsg tbMsg, DeviceToDeviceActorMsg src, FromDeviceRequestMsg fromDeviceRequestMsg) { |
|
|
|
SessionMsgType sessionMsgType = fromDeviceRequestMsg.getMsgType(); |
|
|
|
int requestId = fromDeviceRequestMsg.getRequestId(); |
|
|
|
if (systemContext.isQueuePersistenceEnabled()) { |
|
|
|
pendingMsgs.put(tbMsg.getId(), new PendingSessionMsgData(src.getSessionId(), src.getServerAddress(), sessionMsgType, requestId)); |
|
|
|
scheduleMsgWithDelay(context, new DeviceActorQueueTimeoutMsg(tbMsg.getId(), systemContext.getQueuePersistenceTimeout()), systemContext.getQueuePersistenceTimeout()); |
|
|
|
} else { |
|
|
|
ToDeviceSessionActorMsg response = new BasicToDeviceSessionActorMsg(BasicStatusCodeResponse.onSuccess(sessionMsgType, requestId), src.getSessionId()); |
|
|
|
sendMsgToSessionActor(response, src.getServerAddress()); |
|
|
|
} |
|
|
|
context.parent().tell(new DeviceActorToRuleEngineMsg(context.self(), tbMsg), context.self()); |
|
|
|
} |
|
|
|
|
|
|
|
void processAttributesUpdate(ActorContext context, DeviceAttributesEventNotificationMsg msg) { |
|
|
|
refreshAttributes(msg); |
|
|
|
if (attributeSubscriptions.size() > 0) { |
|
|
|
ToDeviceMsg notification = null; |
|
|
|
if (msg.isDeleted()) { |
|
|
|
@ -225,50 +327,29 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
// void process(ActorContext context, RuleChainDeviceMsg srcMsg) {
|
|
|
|
// ChainProcessingMetaData md = new ChainProcessingMetaData(srcMsg.getRuleChain(),
|
|
|
|
// srcMsg.getToDeviceActorMsg(), new DeviceMetaData(deviceId, deviceName, deviceType, deviceAttributes), context.self());
|
|
|
|
// ChainProcessingContext ctx = new ChainProcessingContext(md);
|
|
|
|
// if (ctx.getChainLength() > 0) {
|
|
|
|
// RuleProcessingMsg msg = new RuleProcessingMsg(ctx);
|
|
|
|
// ActorRef ruleActorRef = ctx.getCurrentActor();
|
|
|
|
// ruleActorRef.tell(msg, ActorRef.noSender());
|
|
|
|
// } else {
|
|
|
|
// context.self().tell(new RulesProcessedMsg(ctx), context.self());
|
|
|
|
// }
|
|
|
|
// }
|
|
|
|
|
|
|
|
void processRpcResponses(ActorContext context, ToDeviceActorMsg msg) { |
|
|
|
private void processRpcResponses(ActorContext context, DeviceToDeviceActorMsg msg) { |
|
|
|
SessionId sessionId = msg.getSessionId(); |
|
|
|
FromDeviceMsg inMsg = msg.getPayload(); |
|
|
|
if (inMsg.getMsgType() == MsgType.TO_DEVICE_RPC_RESPONSE) { |
|
|
|
if (inMsg.getMsgType() == SessionMsgType.TO_DEVICE_RPC_RESPONSE) { |
|
|
|
logger.debug("[{}] Processing rpc command response [{}]", deviceId, sessionId); |
|
|
|
ToDeviceRpcResponseMsg responseMsg = (ToDeviceRpcResponseMsg) inMsg; |
|
|
|
ToDeviceRpcRequestMetadata requestMd = rpcPendingMap.remove(responseMsg.getRequestId()); |
|
|
|
boolean success = requestMd != null; |
|
|
|
if (success) { |
|
|
|
ToPluginRpcResponseDeviceMsg responsePluginMsg = toPluginRpcResponseMsg(requestMd.getMsg(), responseMsg.getData()); |
|
|
|
Optional<ServerAddress> pluginServerAddress = requestMd.getMsg().getServerAddress(); |
|
|
|
if (pluginServerAddress.isPresent()) { |
|
|
|
systemContext.getRpcService().tell(pluginServerAddress.get(), responsePluginMsg); |
|
|
|
logger.debug("[{}] Rpc command response sent to remote plugin actor [{}]!", deviceId, requestMd.getMsg().getMsg().getId()); |
|
|
|
} else { |
|
|
|
context.parent().tell(responsePluginMsg, ActorRef.noSender()); |
|
|
|
logger.debug("[{}] Rpc command response sent to local plugin actor [{}]!", deviceId, requestMd.getMsg().getMsg().getId()); |
|
|
|
} |
|
|
|
systemContext.getDeviceRpcService().process(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), responseMsg.getData(), null)); |
|
|
|
} else { |
|
|
|
logger.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); |
|
|
|
} |
|
|
|
if (msg.getSessionType() == SessionType.SYNC) { |
|
|
|
BasicCommandAckResponse response = success |
|
|
|
? BasicCommandAckResponse.onSuccess(MsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId()) |
|
|
|
: BasicCommandAckResponse.onError(MsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId(), new TimeoutException()); |
|
|
|
? BasicCommandAckResponse.onSuccess(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId()) |
|
|
|
: BasicCommandAckResponse.onError(SessionMsgType.TO_DEVICE_RPC_REQUEST, responseMsg.getRequestId(), new TimeoutException()); |
|
|
|
sendMsgToSessionActor(new BasicToDeviceSessionActorMsg(response, msg.getSessionId()), msg.getServerAddress()); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
public void processClusterEventMsg(ClusterEventMsg msg) { |
|
|
|
void processClusterEventMsg(ClusterEventMsg msg) { |
|
|
|
if (!msg.isAdded()) { |
|
|
|
logger.debug("[{}] Clearing attributes/rpc subscription for server [{}]", deviceId, msg.getServerAddress()); |
|
|
|
Predicate<Map.Entry<SessionId, SessionInfo>> filter = e -> e.getValue().getServer() |
|
|
|
@ -278,59 +359,27 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private ToPluginRpcResponseDeviceMsg toPluginRpcResponseMsg(ToDeviceRpcRequestPluginMsg requestMsg, String data) { |
|
|
|
return toPluginRpcResponseMsg(requestMsg, data, null); |
|
|
|
} |
|
|
|
|
|
|
|
private ToPluginRpcResponseDeviceMsg toPluginRpcResponseMsg(ToDeviceRpcRequestPluginMsg requestMsg, RpcError error) { |
|
|
|
return toPluginRpcResponseMsg(requestMsg, null, error); |
|
|
|
} |
|
|
|
|
|
|
|
private ToPluginRpcResponseDeviceMsg toPluginRpcResponseMsg(ToDeviceRpcRequestPluginMsg requestMsg, String data, RpcError error) { |
|
|
|
return new ToPluginRpcResponseDeviceMsg( |
|
|
|
requestMsg.getPluginId(), |
|
|
|
requestMsg.getPluginTenantId(), |
|
|
|
new FromDeviceRpcResponse(requestMsg.getMsg().getId(), |
|
|
|
data, |
|
|
|
error |
|
|
|
) |
|
|
|
); |
|
|
|
} |
|
|
|
|
|
|
|
// void onRulesProcessedMsg(ActorContext context, RulesProcessedMsg msg) {
|
|
|
|
// ChainProcessingContext ctx = msg.getCtx();
|
|
|
|
// ToDeviceActorMsg inMsg = ctx.getInMsg();
|
|
|
|
// SessionId sid = inMsg.getSessionId();
|
|
|
|
// ToDeviceSessionActorMsg response;
|
|
|
|
// if (ctx.getResponse() != null) {
|
|
|
|
// response = new BasicToDeviceSessionActorMsg(ctx.getResponse(), sid);
|
|
|
|
// } else {
|
|
|
|
// response = new BasicToDeviceSessionActorMsg(ctx.getError(), sid);
|
|
|
|
// }
|
|
|
|
// sendMsgToSessionActor(response, inMsg.getServerAddress());
|
|
|
|
// }
|
|
|
|
|
|
|
|
private void processSubscriptionCommands(ActorContext context, ToDeviceActorMsg msg) { |
|
|
|
private void processSubscriptionCommands(ActorContext context, DeviceToDeviceActorMsg msg) { |
|
|
|
SessionId sessionId = msg.getSessionId(); |
|
|
|
SessionType sessionType = msg.getSessionType(); |
|
|
|
FromDeviceMsg inMsg = msg.getPayload(); |
|
|
|
if (inMsg.getMsgType() == MsgType.SUBSCRIBE_ATTRIBUTES_REQUEST) { |
|
|
|
if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_ATTRIBUTES_REQUEST) { |
|
|
|
logger.debug("[{}] Registering attributes subscription for session [{}]", deviceId, sessionId); |
|
|
|
attributeSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress())); |
|
|
|
} else if (inMsg.getMsgType() == MsgType.UNSUBSCRIBE_ATTRIBUTES_REQUEST) { |
|
|
|
} else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_ATTRIBUTES_REQUEST) { |
|
|
|
logger.debug("[{}] Canceling attributes subscription for session [{}]", deviceId, sessionId); |
|
|
|
attributeSubscriptions.remove(sessionId); |
|
|
|
} else if (inMsg.getMsgType() == MsgType.SUBSCRIBE_RPC_COMMANDS_REQUEST) { |
|
|
|
} else if (inMsg.getMsgType() == SessionMsgType.SUBSCRIBE_RPC_COMMANDS_REQUEST) { |
|
|
|
logger.debug("[{}] Registering rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType); |
|
|
|
rpcSubscriptions.put(sessionId, new SessionInfo(sessionType, msg.getServerAddress())); |
|
|
|
sendPendingRequests(context, sessionId, sessionType, msg.getServerAddress()); |
|
|
|
} else if (inMsg.getMsgType() == MsgType.UNSUBSCRIBE_RPC_COMMANDS_REQUEST) { |
|
|
|
} else if (inMsg.getMsgType() == SessionMsgType.UNSUBSCRIBE_RPC_COMMANDS_REQUEST) { |
|
|
|
logger.debug("[{}] Canceling rpc subscription for session [{}][{}]", deviceId, sessionId, sessionType); |
|
|
|
rpcSubscriptions.remove(sessionId); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void processSessionStateMsgs(ToDeviceActorMsg msg) { |
|
|
|
private void processSessionStateMsgs(DeviceToDeviceActorMsg msg) { |
|
|
|
SessionId sessionId = msg.getSessionId(); |
|
|
|
FromDeviceMsg inMsg = msg.getPayload(); |
|
|
|
if (inMsg instanceof SessionOpenMsg) { |
|
|
|
@ -364,7 +413,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
public void processCredentialsUpdate() { |
|
|
|
void processCredentialsUpdate() { |
|
|
|
sessions.forEach((k, v) -> { |
|
|
|
sendMsgToSessionActor(new BasicToDeviceSessionActorMsg(new SessionCloseNotification(), k), v.getServer()); |
|
|
|
}); |
|
|
|
@ -372,8 +421,12 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso |
|
|
|
rpcSubscriptions.clear(); |
|
|
|
} |
|
|
|
|
|
|
|
public void processNameOrTypeUpdate(DeviceNameOrTypeUpdateMsg msg) { |
|
|
|
void processNameOrTypeUpdate(DeviceNameOrTypeUpdateMsg msg) { |
|
|
|
this.deviceName = msg.getDeviceName(); |
|
|
|
this.deviceType = msg.getDeviceType(); |
|
|
|
this.defaultMetaData = new TbMsgMetaData(); |
|
|
|
this.defaultMetaData.putValue("deviceName", deviceName); |
|
|
|
this.defaultMetaData.putValue("deviceType", deviceType); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|