|
|
|
@ -16,10 +16,6 @@ |
|
|
|
package org.thingsboard.server.transport.coap; |
|
|
|
|
|
|
|
import com.google.gson.JsonParseException; |
|
|
|
import com.google.protobuf.Descriptors; |
|
|
|
import com.google.protobuf.DynamicMessage; |
|
|
|
import lombok.Data; |
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.eclipse.californium.core.coap.CoAP; |
|
|
|
import org.eclipse.californium.core.coap.Request; |
|
|
|
@ -29,7 +25,6 @@ import org.eclipse.californium.core.observe.ObserveRelation; |
|
|
|
import org.eclipse.californium.core.server.resources.CoapExchange; |
|
|
|
import org.eclipse.californium.core.server.resources.Resource; |
|
|
|
import org.eclipse.californium.core.server.resources.ResourceObserver; |
|
|
|
import org.springframework.util.CollectionUtils; |
|
|
|
import org.thingsboard.server.coapserver.CoapServerService; |
|
|
|
import org.thingsboard.server.coapserver.TbCoapDtlsSessionInfo; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
@ -37,38 +32,30 @@ import org.thingsboard.server.common.data.DeviceProfile; |
|
|
|
import org.thingsboard.server.common.data.DeviceTransportType; |
|
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
import org.thingsboard.server.common.data.TransportPayloadType; |
|
|
|
import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.DefaultCoapDeviceTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.JsonTransportPayloadConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.security.DeviceTokenCredentials; |
|
|
|
import org.thingsboard.server.common.msg.session.FeatureType; |
|
|
|
import org.thingsboard.server.common.msg.session.SessionMsgType; |
|
|
|
import org.thingsboard.server.common.transport.SessionMsgListener; |
|
|
|
import org.thingsboard.server.common.transport.TransportServiceCallback; |
|
|
|
import org.thingsboard.server.common.transport.adaptor.AdaptorException; |
|
|
|
import org.thingsboard.server.common.transport.adaptor.JsonConverter; |
|
|
|
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
|
import org.thingsboard.server.transport.coap.adaptors.CoapTransportAdaptor; |
|
|
|
import org.thingsboard.server.transport.coap.callback.CoapDeviceAuthCallback; |
|
|
|
import org.thingsboard.server.transport.coap.callback.CoapNoOpCallback; |
|
|
|
import org.thingsboard.server.transport.coap.callback.CoapOkCallback; |
|
|
|
import org.thingsboard.server.transport.coap.callback.GetAttributesSyncSessionCallback; |
|
|
|
import org.thingsboard.server.transport.coap.callback.ToServerRpcSyncSessionCallback; |
|
|
|
import org.thingsboard.server.transport.coap.client.CoapClientContext; |
|
|
|
import org.thingsboard.server.transport.coap.client.TbCoapClientState; |
|
|
|
|
|
|
|
import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.Random; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.atomic.AtomicInteger; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
@ -80,13 +67,11 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
private static final int REQUEST_ID_POSITION_CERTIFICATE_REQUEST = 4; |
|
|
|
private static final String DTLS_SESSION_ID_KEY = "DTLS_SESSION_ID"; |
|
|
|
|
|
|
|
private final ConcurrentMap<String, CoapObserveSessionInfo> tokenToCoapSessionInfoMap = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<CoapObserveSessionInfo, ObserveRelation> sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); |
|
|
|
private final Set<UUID> rpcSubscriptions = ConcurrentHashMap.newKeySet(); |
|
|
|
private final Set<UUID> attributeSubscriptions = ConcurrentHashMap.newKeySet(); |
|
|
|
private final ConcurrentMap<TbCoapClientState, ObserveRelation> sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private final ConcurrentMap<String, TbCoapDtlsSessionInfo> dtlsSessionIdMap; |
|
|
|
private final long timeout; |
|
|
|
private final CoapClientContext clients; |
|
|
|
|
|
|
|
public CoapTransportResource(CoapTransportContext ctx, CoapServerService coapServerService, String name) { |
|
|
|
super(ctx, name); |
|
|
|
@ -94,17 +79,14 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
this.addObserver(new CoapResourceObserver()); |
|
|
|
this.dtlsSessionIdMap = coapServerService.getDtlsSessionsMap(); |
|
|
|
this.timeout = coapServerService.getTimeout(); |
|
|
|
this.clients = ctx.getClientContext(); |
|
|
|
long sessionReportTimeout = ctx.getSessionReportTimeout(); |
|
|
|
ctx.getScheduler().scheduleAtFixedRate(() -> { |
|
|
|
Set<CoapObserveSessionInfo> coapObserveSessionInfos = sessionInfoToObserveRelationMap.keySet(); |
|
|
|
Set<TransportProtos.SessionInfoProto> observeSessions = coapObserveSessionInfos |
|
|
|
.stream() |
|
|
|
.map(CoapObserveSessionInfo::getSessionInfoProto) |
|
|
|
.collect(Collectors.toSet()); |
|
|
|
observeSessions.forEach(this::reportActivity); |
|
|
|
}, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); |
|
|
|
ctx.getScheduler().scheduleAtFixedRate(clients::reportActivity, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); |
|
|
|
} |
|
|
|
|
|
|
|
/* |
|
|
|
* Overwritten method from CoapResource to be able to manage our own observe notification counters. |
|
|
|
*/ |
|
|
|
@Override |
|
|
|
public void checkObserveRelation(Exchange exchange, Response response) { |
|
|
|
String token = getTokenFromRequest(exchange.getRequest()); |
|
|
|
@ -117,20 +99,15 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
relation.setEstablished(); |
|
|
|
addObserveRelation(relation); |
|
|
|
} |
|
|
|
AtomicInteger observeNotificationCounter = tokenToCoapSessionInfoMap.get(token).getObserveNotificationCounter(); |
|
|
|
response.getOptions().setObserve(observeNotificationCounter.getAndIncrement()); |
|
|
|
AtomicInteger state = clients.getNotificationCounterByToken(token); |
|
|
|
if (state != null) { |
|
|
|
response.getOptions().setObserve(state.getAndIncrement()); |
|
|
|
} else { |
|
|
|
response.getOptions().removeObserve(); |
|
|
|
} |
|
|
|
} // ObserveLayer takes care of the else case
|
|
|
|
} |
|
|
|
|
|
|
|
private void clearAndNotifyObserveRelation(ObserveRelation relation, CoAP.ResponseCode code) { |
|
|
|
relation.cancel(); |
|
|
|
relation.getExchange().sendResponse(new Response(code)); |
|
|
|
} |
|
|
|
|
|
|
|
private Map<CoapObserveSessionInfo, ObserveRelation> getCoapSessionInfoToObserveRelationMap() { |
|
|
|
return sessionInfoToObserveRelationMap; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
protected void processHandleGet(CoapExchange exchange) { |
|
|
|
Optional<FeatureType> featureType = getFeatureType(exchange.advanced().getRequest()); |
|
|
|
@ -232,7 +209,7 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
return dtlsSessionInfo; |
|
|
|
}); |
|
|
|
if (tbCoapDtlsSessionInfo != null) { |
|
|
|
processRequest(exchange, type, request, tbCoapDtlsSessionInfo.getSessionInfoProto(), tbCoapDtlsSessionInfo.getDeviceProfile()); |
|
|
|
processRequest(exchange, type, request, tbCoapDtlsSessionInfo.getMsg(), tbCoapDtlsSessionInfo.getDeviceProfile()); |
|
|
|
} else { |
|
|
|
processAccessTokenRequest(exchange, type, request); |
|
|
|
} |
|
|
|
@ -248,129 +225,133 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
return; |
|
|
|
} |
|
|
|
transportService.process(DeviceTransportType.COAP, TransportProtos.ValidateDeviceTokenRequestMsg.newBuilder().setToken(credentials.get().getCredentialsId()).build(), |
|
|
|
new CoapDeviceAuthCallback(transportContext, exchange, (sessionInfo, deviceProfile) -> { |
|
|
|
processRequest(exchange, type, request, sessionInfo, deviceProfile); |
|
|
|
new CoapDeviceAuthCallback(exchange, (deviceCredentials, deviceProfile) -> { |
|
|
|
processRequest(exchange, type, request, deviceCredentials, deviceProfile); |
|
|
|
})); |
|
|
|
} |
|
|
|
|
|
|
|
private void processRequest(CoapExchange exchange, SessionMsgType type, Request request, TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile) { |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
private void processRequest(CoapExchange exchange, SessionMsgType type, Request request, ValidateDeviceCredentialsResponse deviceCredentials, DeviceProfile deviceProfile) { |
|
|
|
TbCoapClientState clientState = null; |
|
|
|
try { |
|
|
|
TransportConfigurationContainer transportConfigurationContainer = getTransportConfigurationContainer(deviceProfile); |
|
|
|
CoapTransportAdaptor coapTransportAdaptor = getCoapTransportAdaptor(transportConfigurationContainer.isJsonPayload()); |
|
|
|
clientState = clients.getOrCreateClient(type, deviceCredentials, deviceProfile); |
|
|
|
clientState.updateLastUplinkTime(); |
|
|
|
switch (type) { |
|
|
|
case POST_ATTRIBUTES_REQUEST: |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToPostAttributes(sessionId, request, |
|
|
|
transportConfigurationContainer.getAttributesMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
reportSubscriptionInfo(sessionInfo, attributeSubscriptions.contains(sessionId), rpcSubscriptions.contains(sessionId)); |
|
|
|
handlePostAttributesRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case POST_TELEMETRY_REQUEST: |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToPostTelemetry(sessionId, request, |
|
|
|
transportConfigurationContainer.getTelemetryMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
reportSubscriptionInfo(sessionInfo, attributeSubscriptions.contains(sessionId), rpcSubscriptions.contains(sessionId)); |
|
|
|
handlePostTelemetryRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case CLAIM_REQUEST: |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToClaimDevice(sessionId, request, sessionInfo), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
handleClaimRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case SUBSCRIBE_ATTRIBUTES_REQUEST: |
|
|
|
CoapObserveSessionInfo currentCoapObserveAttrSessionInfo = tokenToCoapSessionInfoMap.get(getTokenFromRequest(request)); |
|
|
|
if (currentCoapObserveAttrSessionInfo == null) { |
|
|
|
attributeSubscriptions.add(sessionId); |
|
|
|
registerAsyncCoapSession(exchange, coapTransportAdaptor, transportConfigurationContainer.getRpcRequestDynamicMessageBuilder(), |
|
|
|
sessionInfo, getTokenFromRequest(request)); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
TransportProtos.SubscribeToAttributeUpdatesMsg.getDefaultInstance(), new CoapNoOpCallback(exchange)); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
TransportProtos.GetAttributeRequestMsg.newBuilder().setOnlyShared(true).build(), |
|
|
|
new CoapNoOpCallback(exchange)); |
|
|
|
} |
|
|
|
handleAttributeSubscribeRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case UNSUBSCRIBE_ATTRIBUTES_REQUEST: |
|
|
|
CoapObserveSessionInfo coapObserveAttrSessionInfo = lookupAsyncSessionInfo(getTokenFromRequest(request)); |
|
|
|
if (coapObserveAttrSessionInfo != null) { |
|
|
|
TransportProtos.SessionInfoProto attrSession = coapObserveAttrSessionInfo.getSessionInfoProto(); |
|
|
|
UUID attrSessionId = toSessionId(attrSession); |
|
|
|
attributeSubscriptions.remove(attrSessionId); |
|
|
|
transportService.process(attrSession, |
|
|
|
TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setUnsubscribe(true).build(), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
closeAndDeregister(sessionInfo); |
|
|
|
handleAttributeUnsubscribeRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case SUBSCRIBE_RPC_COMMANDS_REQUEST: |
|
|
|
CoapObserveSessionInfo currentCoapObserveRpcSessionInfo = tokenToCoapSessionInfoMap.get(getTokenFromRequest(request)); |
|
|
|
if (currentCoapObserveRpcSessionInfo == null) { |
|
|
|
rpcSubscriptions.add(sessionId); |
|
|
|
registerAsyncCoapSession(exchange, coapTransportAdaptor, transportConfigurationContainer.getRpcRequestDynamicMessageBuilder() |
|
|
|
, sessionInfo, getTokenFromRequest(request)); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
TransportProtos.SubscribeToRPCMsg.getDefaultInstance(), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.VALID, CoAP.ResponseCode.INTERNAL_SERVER_ERROR) |
|
|
|
); |
|
|
|
} |
|
|
|
handleRpcSubscribeRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case UNSUBSCRIBE_RPC_COMMANDS_REQUEST: |
|
|
|
CoapObserveSessionInfo coapObserveRpcSessionInfo = lookupAsyncSessionInfo(getTokenFromRequest(request)); |
|
|
|
if (coapObserveRpcSessionInfo != null) { |
|
|
|
TransportProtos.SessionInfoProto rpcSession = coapObserveRpcSessionInfo.getSessionInfoProto(); |
|
|
|
UUID rpcSessionId = toSessionId(rpcSession); |
|
|
|
rpcSubscriptions.remove(rpcSessionId); |
|
|
|
transportService.process(rpcSession, |
|
|
|
TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
closeAndDeregister(sessionInfo); |
|
|
|
handleRpcUnsubscribeRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case TO_DEVICE_RPC_RESPONSE: |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToDeviceRpcResponse(sessionId, request, transportConfigurationContainer.getRpcResponseMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
handleToDeviceRpcResponse(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case TO_SERVER_RPC_REQUEST: |
|
|
|
transportService.registerSyncSession(sessionInfo, getCoapSessionListener(exchange, coapTransportAdaptor, |
|
|
|
transportConfigurationContainer.getRpcRequestDynamicMessageBuilder(), sessionInfo), timeout); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToServerRpcRequest(sessionId, request), |
|
|
|
new CoapNoOpCallback(exchange)); |
|
|
|
handleToServerRpcRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
case GET_ATTRIBUTES_REQUEST: |
|
|
|
transportService.registerSyncSession(sessionInfo, getCoapSessionListener(exchange, coapTransportAdaptor, |
|
|
|
transportConfigurationContainer.getRpcRequestDynamicMessageBuilder(), sessionInfo), timeout); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
coapTransportAdaptor.convertToGetAttributes(sessionId, request), |
|
|
|
new CoapNoOpCallback(exchange)); |
|
|
|
handleGetAttributesRequest(clientState, exchange, request); |
|
|
|
break; |
|
|
|
} |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.trace("[{}] Failed to decode message: ", sessionId, e); |
|
|
|
if (clientState != null) { |
|
|
|
log.trace("[{}] Failed to decode message: ", clientState.getDeviceId(), e); |
|
|
|
} |
|
|
|
exchange.respond(CoAP.ResponseCode.BAD_REQUEST); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private UUID toSessionId(TransportProtos.SessionInfoProto sessionInfoProto) { |
|
|
|
return new UUID(sessionInfoProto.getSessionIdMSB(), sessionInfoProto.getSessionIdLSB()); |
|
|
|
private void handlePostAttributesRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = clients.getNewSyncSession(clientState); |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
transportService.process(sessionInfo, clientState.getAdaptor().convertToPostAttributes(sessionId, request, |
|
|
|
clientState.getConfiguration().getAttributesMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
|
|
|
|
private void handlePostTelemetryRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = clients.getNewSyncSession(clientState); |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
transportService.process(sessionInfo, clientState.getAdaptor().convertToPostTelemetry(sessionId, request, |
|
|
|
clientState.getConfiguration().getTelemetryMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
|
|
|
|
private void handleClaimRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = clients.getNewSyncSession(clientState); |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
clientState.getAdaptor().convertToClaimDevice(sessionId, request, sessionInfo), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
|
|
|
|
private void handleAttributeSubscribeRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) { |
|
|
|
String attrSubToken = getTokenFromRequest(request); |
|
|
|
if (!clients.registerAttributeObservation(clientState, attrSubToken, exchange)) { |
|
|
|
log.warn("[{}] Received duplicate attribute subscribe request for token: {}", clientState.getDeviceId(), attrSubToken); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void handleAttributeUnsubscribeRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) { |
|
|
|
clients.deregisterAttributeObservation(clientState, getTokenFromRequest(request), exchange); |
|
|
|
} |
|
|
|
|
|
|
|
private CoapObserveSessionInfo lookupAsyncSessionInfo(String token) { |
|
|
|
return tokenToCoapSessionInfoMap.remove(token); |
|
|
|
private void handleRpcUnsubscribeRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) { |
|
|
|
clients.deregisterRpcObservation(clientState, getTokenFromRequest(request), exchange); |
|
|
|
} |
|
|
|
|
|
|
|
private void registerAsyncCoapSession(CoapExchange exchange, CoapTransportAdaptor coapTransportAdaptor, |
|
|
|
DynamicMessage.Builder rpcRequestDynamicMessageBuilder, TransportProtos.SessionInfoProto sessionInfo, String token) { |
|
|
|
tokenToCoapSessionInfoMap.putIfAbsent(token, new CoapObserveSessionInfo(sessionInfo)); |
|
|
|
transportService.registerAsyncSession(sessionInfo, getCoapSessionListener(exchange, coapTransportAdaptor, rpcRequestDynamicMessageBuilder, sessionInfo)); |
|
|
|
transportService.process(sessionInfo, getSessionEventMsg(TransportProtos.SessionEvent.OPEN), null); |
|
|
|
private void handleToDeviceRpcResponse(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto session = clientState.getSession(); |
|
|
|
if (session == null) { |
|
|
|
session = clients.getNewSyncSession(clientState); |
|
|
|
} |
|
|
|
UUID sessionId = toSessionId(session); |
|
|
|
transportService.process(session, |
|
|
|
clientState.getAdaptor().convertToDeviceRpcResponse(sessionId, request, clientState.getConfiguration().getRpcResponseMsgDescriptor()), |
|
|
|
new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); |
|
|
|
} |
|
|
|
|
|
|
|
private CoapSessionListener getCoapSessionListener(CoapExchange exchange, CoapTransportAdaptor coapTransportAdaptor, |
|
|
|
DynamicMessage.Builder rpcRequestDynamicMessageBuilder, TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
return new CoapSessionListener(exchange, coapTransportAdaptor, rpcRequestDynamicMessageBuilder, sessionInfo); |
|
|
|
private void handleRpcSubscribeRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) { |
|
|
|
String rpcSubToken = getTokenFromRequest(request); |
|
|
|
if (!clients.registerRpcObservation(clientState, rpcSubToken, exchange)) { |
|
|
|
log.warn("[{}] Received duplicate rpc subscribe request.", rpcSubToken); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void handleGetAttributesRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = clients.getNewSyncSession(clientState); |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
transportService.registerSyncSession(sessionInfo, new GetAttributesSyncSessionCallback(clientState, exchange, request), timeout); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
clientState.getAdaptor().convertToGetAttributes(sessionId, request), |
|
|
|
new CoapNoOpCallback(exchange)); |
|
|
|
} |
|
|
|
|
|
|
|
private void handleToServerRpcRequest(TbCoapClientState clientState, CoapExchange exchange, Request request) throws AdaptorException { |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = clients.getNewSyncSession(clientState); |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
transportService.registerSyncSession(sessionInfo, new ToServerRpcSyncSessionCallback(clientState, exchange, request), timeout); |
|
|
|
transportService.process(sessionInfo, |
|
|
|
clientState.getAdaptor().convertToServerRpcRequest(sessionId, request), |
|
|
|
new CoapNoOpCallback(exchange)); |
|
|
|
} |
|
|
|
|
|
|
|
private UUID toSessionId(TransportProtos.SessionInfoProto sessionInfoProto) { |
|
|
|
return new UUID(sessionInfoProto.getSessionIdMSB(), sessionInfoProto.getSessionIdLSB()); |
|
|
|
} |
|
|
|
|
|
|
|
private String getTokenFromRequest(Request request) { |
|
|
|
@ -452,119 +433,6 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@RequiredArgsConstructor |
|
|
|
private class CoapSessionListener implements SessionMsgListener { |
|
|
|
|
|
|
|
private final CoapExchange exchange; |
|
|
|
private final CoapTransportAdaptor coapTransportAdaptor; |
|
|
|
private final DynamicMessage.Builder rpcRequestDynamicMessageBuilder; |
|
|
|
private final TransportProtos.SessionInfoProto sessionInfo; |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg msg) { |
|
|
|
try { |
|
|
|
exchange.respond(coapTransportAdaptor.convertToPublish(isConRequest(), msg)); |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.trace("Failed to reply due to error", e); |
|
|
|
exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg msg) { |
|
|
|
log.trace("[{}] Received attributes update notification to device", sessionId); |
|
|
|
try { |
|
|
|
exchange.respond(coapTransportAdaptor.convertToPublish(isConRequest(), msg)); |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.trace("Failed to reply due to error", e); |
|
|
|
closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.INTERNAL_SERVER_ERROR); |
|
|
|
closeAndDeregister(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { |
|
|
|
log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); |
|
|
|
closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.SERVICE_UNAVAILABLE); |
|
|
|
closeAndDeregister(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg msg) { |
|
|
|
log.trace("[{}] Received RPC command to device", sessionId); |
|
|
|
boolean sent = false; |
|
|
|
try { |
|
|
|
Response response = coapTransportAdaptor.convertToPublish(isConRequest(), msg, rpcRequestDynamicMessageBuilder); |
|
|
|
int requestId = getNextMsgId(); |
|
|
|
response.setMID(requestId); |
|
|
|
|
|
|
|
if (msg.getPersisted() && isConRequest()) { |
|
|
|
transportContext.getRpcAwaitingAck().put(requestId, msg); |
|
|
|
transportContext.getScheduler().schedule(() -> { |
|
|
|
TransportProtos.ToDeviceRpcRequestMsg awaitingAckMsg = transportContext.getRpcAwaitingAck().remove(requestId); |
|
|
|
if (awaitingAckMsg != null) { |
|
|
|
transportService.process(sessionInfo, msg, true, TransportServiceCallback.EMPTY); |
|
|
|
} |
|
|
|
}, Math.max(0, msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); |
|
|
|
response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { |
|
|
|
TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); |
|
|
|
if (rpcRequestMsg != null) { |
|
|
|
transportService.process(sessionInfo, rpcRequestMsg, false, TransportServiceCallback.EMPTY); |
|
|
|
} |
|
|
|
})); |
|
|
|
} |
|
|
|
exchange.respond(response); |
|
|
|
sent = true; |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.trace("Failed to reply due to error", e); |
|
|
|
closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.INTERNAL_SERVER_ERROR); |
|
|
|
closeAndDeregister(); |
|
|
|
} finally { |
|
|
|
if (msg.getPersisted() && !isConRequest()) { |
|
|
|
transportService.process(sessionInfo, msg, sent, TransportServiceCallback.EMPTY); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg msg) { |
|
|
|
try { |
|
|
|
exchange.respond(coapTransportAdaptor.convertToPublish(isConRequest(), msg)); |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.trace("Failed to reply due to error", e); |
|
|
|
exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private boolean isConRequest() { |
|
|
|
return exchange.advanced().getRequest().isConfirmable(); |
|
|
|
} |
|
|
|
|
|
|
|
private void closeObserveRelationAndNotify(UUID sessionId, CoAP.ResponseCode responseCode) { |
|
|
|
Map<CoapObserveSessionInfo, ObserveRelation> sessionToObserveRelationMap = CoapTransportResource.this.getCoapSessionInfoToObserveRelationMap(); |
|
|
|
if (CoapTransportResource.this.getObserverCount() > 0 && !CollectionUtils.isEmpty(sessionToObserveRelationMap)) { |
|
|
|
Optional<CoapObserveSessionInfo> observeSessionToClose = sessionToObserveRelationMap.keySet().stream().filter(coapObserveSessionInfo -> { |
|
|
|
TransportProtos.SessionInfoProto sessionToDelete = coapObserveSessionInfo.getSessionInfoProto(); |
|
|
|
UUID observeSessionId = new UUID(sessionToDelete.getSessionIdMSB(), sessionToDelete.getSessionIdLSB()); |
|
|
|
return observeSessionId.equals(sessionId); |
|
|
|
}).findFirst(); |
|
|
|
if (observeSessionToClose.isPresent()) { |
|
|
|
CoapObserveSessionInfo coapObserveSessionInfo = observeSessionToClose.get(); |
|
|
|
ObserveRelation observeRelation = sessionToObserveRelationMap.get(coapObserveSessionInfo); |
|
|
|
CoapTransportResource.this.clearAndNotifyObserveRelation(observeRelation, responseCode); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void closeAndDeregister() { |
|
|
|
Request request = exchange.advanced().getRequest(); |
|
|
|
String token = CoapTransportResource.this.getTokenFromRequest(request); |
|
|
|
CoapObserveSessionInfo deleted = CoapTransportResource.this.lookupAsyncSessionInfo(token); |
|
|
|
CoapTransportResource.this.closeAndDeregister(deleted.getSessionInfoProto()); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
public class CoapResourceObserver implements ResourceObserver { |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -587,7 +455,7 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
public void addedObserveRelation(ObserveRelation relation) { |
|
|
|
Request request = relation.getExchange().getRequest(); |
|
|
|
String token = getTokenFromRequest(request); |
|
|
|
sessionInfoToObserveRelationMap.putIfAbsent(tokenToCoapSessionInfoMap.get(token), relation); |
|
|
|
clients.registerObserveRelation(token, relation); |
|
|
|
log.trace("Added Observe relation for token: {}", token); |
|
|
|
} |
|
|
|
|
|
|
|
@ -595,93 +463,10 @@ public class CoapTransportResource extends AbstractCoapTransportResource { |
|
|
|
public void removedObserveRelation(ObserveRelation relation) { |
|
|
|
Request request = relation.getExchange().getRequest(); |
|
|
|
String token = getTokenFromRequest(request); |
|
|
|
sessionInfoToObserveRelationMap.remove(tokenToCoapSessionInfoMap.get(token)); |
|
|
|
clients.deregisterObserveRelation(token); |
|
|
|
log.trace("Relation removed for token: {}", token); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void closeAndDeregister(TransportProtos.SessionInfoProto session) { |
|
|
|
UUID sessionId = toSessionId(session); |
|
|
|
transportService.process(session, getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); |
|
|
|
transportService.deregisterSession(session); |
|
|
|
rpcSubscriptions.remove(sessionId); |
|
|
|
attributeSubscriptions.remove(sessionId); |
|
|
|
} |
|
|
|
|
|
|
|
private TransportConfigurationContainer getTransportConfigurationContainer(DeviceProfile deviceProfile) throws AdaptorException { |
|
|
|
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); |
|
|
|
if (transportConfiguration instanceof DefaultDeviceProfileTransportConfiguration) { |
|
|
|
return new TransportConfigurationContainer(true); |
|
|
|
} else if (transportConfiguration instanceof CoapDeviceProfileTransportConfiguration) { |
|
|
|
CoapDeviceProfileTransportConfiguration coapDeviceProfileTransportConfiguration = |
|
|
|
(CoapDeviceProfileTransportConfiguration) transportConfiguration; |
|
|
|
CoapDeviceTypeConfiguration coapDeviceTypeConfiguration = |
|
|
|
coapDeviceProfileTransportConfiguration.getCoapDeviceTypeConfiguration(); |
|
|
|
if (coapDeviceTypeConfiguration instanceof DefaultCoapDeviceTypeConfiguration) { |
|
|
|
DefaultCoapDeviceTypeConfiguration defaultCoapDeviceTypeConfiguration = |
|
|
|
(DefaultCoapDeviceTypeConfiguration) coapDeviceTypeConfiguration; |
|
|
|
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = |
|
|
|
defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); |
|
|
|
if (transportPayloadTypeConfiguration instanceof JsonTransportPayloadConfiguration) { |
|
|
|
return new TransportConfigurationContainer(true); |
|
|
|
} else { |
|
|
|
ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = |
|
|
|
(ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; |
|
|
|
String deviceTelemetryProtoSchema = protoTransportPayloadConfiguration.getDeviceTelemetryProtoSchema(); |
|
|
|
String deviceAttributesProtoSchema = protoTransportPayloadConfiguration.getDeviceAttributesProtoSchema(); |
|
|
|
String deviceRpcRequestProtoSchema = protoTransportPayloadConfiguration.getDeviceRpcRequestProtoSchema(); |
|
|
|
String deviceRpcResponseProtoSchema = protoTransportPayloadConfiguration.getDeviceRpcResponseProtoSchema(); |
|
|
|
return new TransportConfigurationContainer(false, |
|
|
|
protoTransportPayloadConfiguration.getTelemetryDynamicMessageDescriptor(deviceTelemetryProtoSchema), |
|
|
|
protoTransportPayloadConfiguration.getAttributesDynamicMessageDescriptor(deviceAttributesProtoSchema), |
|
|
|
protoTransportPayloadConfiguration.getRpcResponseDynamicMessageDescriptor(deviceRpcResponseProtoSchema), |
|
|
|
protoTransportPayloadConfiguration.getRpcRequestDynamicMessageBuilder(deviceRpcRequestProtoSchema) |
|
|
|
); |
|
|
|
} |
|
|
|
} else { |
|
|
|
throw new AdaptorException("Invalid CoapDeviceTypeConfiguration type: " + coapDeviceTypeConfiguration.getClass().getSimpleName() + "!"); |
|
|
|
} |
|
|
|
} else { |
|
|
|
throw new AdaptorException("Invalid DeviceProfileTransportConfiguration type" + transportConfiguration.getClass().getSimpleName() + "!"); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private CoapTransportAdaptor getCoapTransportAdaptor(boolean jsonPayloadType) { |
|
|
|
return jsonPayloadType ? transportContext.getJsonCoapAdaptor() : transportContext.getProtoCoapAdaptor(); |
|
|
|
} |
|
|
|
|
|
|
|
@Data |
|
|
|
private static class TransportConfigurationContainer { |
|
|
|
|
|
|
|
private boolean jsonPayload; |
|
|
|
private Descriptors.Descriptor telemetryMsgDescriptor; |
|
|
|
private Descriptors.Descriptor attributesMsgDescriptor; |
|
|
|
private Descriptors.Descriptor rpcResponseMsgDescriptor; |
|
|
|
private DynamicMessage.Builder rpcRequestDynamicMessageBuilder; |
|
|
|
|
|
|
|
public TransportConfigurationContainer(boolean jsonPayload, Descriptors.Descriptor telemetryMsgDescriptor, Descriptors.Descriptor attributesMsgDescriptor, Descriptors.Descriptor rpcResponseMsgDescriptor, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) { |
|
|
|
this.jsonPayload = jsonPayload; |
|
|
|
this.telemetryMsgDescriptor = telemetryMsgDescriptor; |
|
|
|
this.attributesMsgDescriptor = attributesMsgDescriptor; |
|
|
|
this.rpcResponseMsgDescriptor = rpcResponseMsgDescriptor; |
|
|
|
this.rpcRequestDynamicMessageBuilder = rpcRequestDynamicMessageBuilder; |
|
|
|
} |
|
|
|
|
|
|
|
public TransportConfigurationContainer(boolean jsonPayload) { |
|
|
|
this.jsonPayload = jsonPayload; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Data |
|
|
|
private static class CoapObserveSessionInfo { |
|
|
|
|
|
|
|
private final TransportProtos.SessionInfoProto sessionInfoProto; |
|
|
|
private final AtomicInteger observeNotificationCounter; |
|
|
|
|
|
|
|
private CoapObserveSessionInfo(TransportProtos.SessionInfoProto sessionInfoProto) { |
|
|
|
this.sessionInfoProto = sessionInfoProto; |
|
|
|
this.observeNotificationCounter = new AtomicInteger(0); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|