|
|
|
@ -17,15 +17,12 @@ package org.thingsboard.server.transport.lwm2m.server.rpc; |
|
|
|
|
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.apache.commons.lang3.exception.ExceptionUtils; |
|
|
|
import org.eclipse.leshan.core.ResponseCode; |
|
|
|
import org.eclipse.leshan.core.request.ReadCompositeRequest; |
|
|
|
import org.eclipse.leshan.core.response.ReadCompositeResponse; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
import org.thingsboard.server.common.data.rpc.RpcStatus; |
|
|
|
import org.thingsboard.server.common.transport.TransportService; |
|
|
|
import org.thingsboard.server.common.transport.TransportServiceCallback; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
|
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; |
|
|
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
|
|
|
@ -63,11 +60,9 @@ import org.thingsboard.server.transport.lwm2m.server.rpc.composite.RpcReadRespon |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.rpc.composite.RpcWriteCompositeRequest; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; |
|
|
|
|
|
|
|
import java.util.Map; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
@Service |
|
|
|
@ -85,91 +80,96 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { |
|
|
|
@Override |
|
|
|
public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg rpcRequest, TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
log.debug("Received params: {}", rpcRequest.getParams()); |
|
|
|
LwM2mOperationType operationType = LwM2mOperationType.fromType(rpcRequest.getMethodName()); |
|
|
|
if (operationType == null) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.METHOD_NOT_ALLOWED, "Unsupported operation type: " + rpcRequest.getMethodName()); |
|
|
|
return; |
|
|
|
} |
|
|
|
LwM2mClient client = clientContext.getClientBySessionInfo(sessionInfo); |
|
|
|
if (client.getRegistration() == null) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.INTERNAL_SERVER_ERROR, "Registration is empty"); |
|
|
|
return; |
|
|
|
} |
|
|
|
UUID rpcId = new UUID(rpcRequest.getRequestIdMSB(), rpcRequest.getRequestIdLSB()); |
|
|
|
|
|
|
|
if (rpcId.equals(client.getLastSentRpcId())) { |
|
|
|
log.debug("[{}]][{}] Rpc has already sent!", client.getEndpoint(), rpcId); |
|
|
|
return; |
|
|
|
} |
|
|
|
try { |
|
|
|
if (operationType.isHasObjectId()) { |
|
|
|
String objectId = getIdFromParameters(client, rpcRequest); |
|
|
|
switch (operationType) { |
|
|
|
case READ: |
|
|
|
sendReadRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case OBSERVE: |
|
|
|
sendObserveRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case DISCOVER: |
|
|
|
sendDiscoverRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case EXECUTE: |
|
|
|
sendExecuteRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_ATTRIBUTES: |
|
|
|
sendWriteAttributesRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case OBSERVE_CANCEL: |
|
|
|
sendCancelObserveRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case DELETE: |
|
|
|
sendDeleteRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_UPDATE: |
|
|
|
sendWriteUpdateRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_REPLACE: |
|
|
|
sendWriteReplaceRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
default: |
|
|
|
throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); |
|
|
|
} |
|
|
|
} else if (operationType.isComposite()) { |
|
|
|
if (clientContext.isComposite(client)) { |
|
|
|
LwM2mOperationType operationType = LwM2mOperationType.fromType(rpcRequest.getMethodName()); |
|
|
|
if (operationType == null) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.METHOD_NOT_ALLOWED, "Unsupported operation type: " + rpcRequest.getMethodName()); |
|
|
|
return; |
|
|
|
} |
|
|
|
LwM2mClient client = clientContext.getClientBySessionInfo(sessionInfo); |
|
|
|
if (client.getRegistration() == null) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.INTERNAL_SERVER_ERROR, "Registration is empty"); |
|
|
|
return; |
|
|
|
} |
|
|
|
UUID rpcId = new UUID(rpcRequest.getRequestIdMSB(), rpcRequest.getRequestIdLSB()); |
|
|
|
|
|
|
|
if (rpcId.equals(client.getLastSentRpcId())) { |
|
|
|
log.debug("[{}]][{}] Rpc has already sent!", client.getEndpoint(), rpcId); |
|
|
|
return; |
|
|
|
} |
|
|
|
try { |
|
|
|
if (operationType.isHasObjectId()) { |
|
|
|
String objectId = getIdFromParameters(client, rpcRequest); |
|
|
|
switch (operationType) { |
|
|
|
case READ_COMPOSITE: |
|
|
|
sendReadCompositeRequest(client, rpcRequest); |
|
|
|
case READ: |
|
|
|
sendReadRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case OBSERVE: |
|
|
|
sendObserveRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case DISCOVER: |
|
|
|
sendDiscoverRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case EXECUTE: |
|
|
|
sendExecuteRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_ATTRIBUTES: |
|
|
|
sendWriteAttributesRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_COMPOSITE: |
|
|
|
sendWriteCompositeRequest(client, rpcRequest); |
|
|
|
case OBSERVE_CANCEL: |
|
|
|
sendCancelObserveRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case DELETE: |
|
|
|
sendDeleteRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_UPDATE: |
|
|
|
sendWriteUpdateRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
case WRITE_REPLACE: |
|
|
|
sendWriteReplaceRequest(client, rpcRequest, objectId); |
|
|
|
break; |
|
|
|
default: |
|
|
|
throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); |
|
|
|
} |
|
|
|
} else if (operationType.isComposite()) { |
|
|
|
if (clientContext.isComposite(client)) { |
|
|
|
switch (operationType) { |
|
|
|
case READ_COMPOSITE: |
|
|
|
sendReadCompositeRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case WRITE_COMPOSITE: |
|
|
|
sendWriteCompositeRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
default: |
|
|
|
throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); |
|
|
|
} |
|
|
|
} else { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), |
|
|
|
ResponseCode.INTERNAL_SERVER_ERROR, "This device does not support Composite Operation"); |
|
|
|
} |
|
|
|
} else { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), |
|
|
|
ResponseCode.INTERNAL_SERVER_ERROR, "This device does not support Composite Operation"); |
|
|
|
} |
|
|
|
} else { |
|
|
|
switch (operationType) { |
|
|
|
case OBSERVE_CANCEL_ALL: |
|
|
|
sendCancelAllObserveRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case OBSERVE_READ_ALL: |
|
|
|
sendObserveAllRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case DISCOVER_ALL: |
|
|
|
sendDiscoverAllRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case FW_UPDATE: |
|
|
|
//TODO: implement and add break statement
|
|
|
|
default: |
|
|
|
throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); |
|
|
|
switch (operationType) { |
|
|
|
case OBSERVE_CANCEL_ALL: |
|
|
|
sendCancelAllObserveRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case OBSERVE_READ_ALL: |
|
|
|
sendObserveAllRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case DISCOVER_ALL: |
|
|
|
sendDiscoverAllRequest(client, rpcRequest); |
|
|
|
break; |
|
|
|
case FW_UPDATE: |
|
|
|
//TODO: implement and add break statement
|
|
|
|
default: |
|
|
|
throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (IllegalArgumentException e) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.BAD_REQUEST, e.getMessage()); |
|
|
|
} |
|
|
|
} catch (IllegalArgumentException e) { |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.BAD_REQUEST, e.getMessage()); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}] Failed to send RPC: [{}]", sessionInfo, rpcRequest, e); |
|
|
|
this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.INTERNAL_SERVER_ERROR, ExceptionUtils.getRootCauseMessage(e)); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|