diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesIntegrationTest.java index 5286accd09..b1a4dc3c10 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesIntegrationTest.java @@ -109,7 +109,7 @@ public abstract class AbstractCoapAttributesUpdatesIntegrationTest extends Abstr protected void validateCurrentStateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(0, callback.getObserve().intValue()); String response = new String(callback.getPayloadBytes(), StandardCharsets.UTF_8); assertEquals(JacksonUtil.toJsonNode(POST_ATTRIBUTES_PAYLOAD_ON_CURRENT_STATE_NOTIFICATION), JacksonUtil.toJsonNode(response)); @@ -118,7 +118,7 @@ public abstract class AbstractCoapAttributesUpdatesIntegrationTest extends Abstr protected void validateEmptyCurrentStateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(0, callback.getObserve().intValue()); String response = new String(callback.getPayloadBytes(), StandardCharsets.UTF_8); assertEquals("{}", response); @@ -127,7 +127,7 @@ public abstract class AbstractCoapAttributesUpdatesIntegrationTest extends Abstr protected void validateUpdateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(1, callback.getObserve().intValue()); String response = new String(callback.getPayloadBytes(), StandardCharsets.UTF_8); assertEquals(JacksonUtil.toJsonNode(POST_ATTRIBUTES_PAYLOAD), JacksonUtil.toJsonNode(response)); @@ -136,7 +136,7 @@ public abstract class AbstractCoapAttributesUpdatesIntegrationTest extends Abstr protected void validateDeleteAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(2, callback.getObserve().intValue()); String response = new String(callback.getPayloadBytes(), StandardCharsets.UTF_8); assertEquals(JacksonUtil.toJsonNode(RESPONSE_ATTRIBUTES_PAYLOAD_DELETED), JacksonUtil.toJsonNode(response)); diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesProtoIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesProtoIntegrationTest.java index 625aaf1592..c3e56ccc62 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesProtoIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/updates/AbstractCoapAttributesUpdatesProtoIntegrationTest.java @@ -62,7 +62,7 @@ public abstract class AbstractCoapAttributesUpdatesProtoIntegrationTest extends protected void validateCurrentStateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(0, callback.getObserve().intValue()); TransportProtos.AttributeUpdateNotificationMsg.Builder expectedCurrentStateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); TransportProtos.TsKvProto tsKvProtoAttribute1 = getTsKvProto("attribute1", "value", TransportProtos.KeyValueType.STRING_V); @@ -90,14 +90,14 @@ public abstract class AbstractCoapAttributesUpdatesProtoIntegrationTest extends protected void validateEmptyCurrentStateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(0, callback.getObserve().intValue()); } protected void validateUpdateAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(1, callback.getObserve().intValue()); TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); List tsKvProtoList = getTsKvProtoList(); @@ -117,7 +117,7 @@ public abstract class AbstractCoapAttributesUpdatesProtoIntegrationTest extends protected void validateDeleteAttributesResponse(TestCoapCallback callback) throws InvalidProtocolBufferException { assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(2, callback.getObserve().intValue()); TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); attributeUpdateNotificationMsgBuilder.addSharedDeleted("attribute5"); diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java index 382f5595a5..e2567576f3 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java @@ -196,7 +196,7 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC assertTrue(StringUtils.isEmpty(result)); assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(1, callback.getObserve().intValue()); } @@ -204,7 +204,7 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC assertEquals(expectedResult, actualResult); assertNotNull(callback.getPayloadBytes()); assertNotNull(callback.getObserve()); - assertEquals(callback.getResponseCode(), CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); assertEquals(expectedObserveNumber, callback.getObserve().intValue()); } diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java index 92fdef3c84..1ddfd0bb63 100644 --- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java +++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java @@ -111,8 +111,7 @@ public class TbCoapDtlsCertificateVerifier implements NewAdvancedCertificateVeri if (msg != null && strCert.equals(msg.getCredentials())) { DeviceProfile deviceProfile = msg.getDeviceProfile(); if (msg.hasDeviceInfo() && deviceProfile != null) { - TransportProtos.SessionInfoProto sessionInfoProto = SessionInfoCreator.create(msg, serviceInfoProvider.getServiceId(), UUID.randomUUID()); - tbCoapDtlsSessionInMemoryStorage.put(session.getSessionIdentifier().toString(), new TbCoapDtlsSessionInfo(sessionInfoProto, deviceProfile)); + tbCoapDtlsSessionInMemoryStorage.put(session.getSessionIdentifier().toString(), new TbCoapDtlsSessionInfo(msg, deviceProfile)); } break; } diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInfo.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInfo.java index 893ca38a5c..3dba715f4f 100644 --- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInfo.java +++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInfo.java @@ -17,18 +17,19 @@ package org.thingsboard.server.coapserver; import lombok.Data; import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.gen.transport.TransportProtos; @Data public class TbCoapDtlsSessionInfo { - private TransportProtos.SessionInfoProto sessionInfoProto; + private ValidateDeviceCredentialsResponse msg; private DeviceProfile deviceProfile; private long lastActivityTime; - public TbCoapDtlsSessionInfo(TransportProtos.SessionInfoProto sessionInfoProto, DeviceProfile deviceProfile) { - this.sessionInfoProto = sessionInfoProto; + public TbCoapDtlsSessionInfo(ValidateDeviceCredentialsResponse msg, DeviceProfile deviceProfile) { + this.msg = msg; this.deviceProfile = deviceProfile; this.lastActivityTime = System.currentTimeMillis(); } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/AbstractCoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/AbstractCoapTransportResource.java index 83a6d2ff11..55b49e9fba 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/AbstractCoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/AbstractCoapTransportResource.java @@ -72,18 +72,4 @@ public abstract class AbstractCoapTransportResource extends CoapResource { .build(), TransportServiceCallback.EMPTY); } - protected void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { - transportService.reportActivity(sessionInfo); - } - - protected static TransportProtos.SessionEventMsg getSessionEventMsg(TransportProtos.SessionEvent event) { - return TransportProtos.SessionEventMsg.newBuilder() - .setSessionType(TransportProtos.SessionType.ASYNC) - .setEvent(event).build(); - } - - protected int getNextMsgId() { - return ThreadLocalRandom.current().nextInt(NONE, MAX_MID + 1); - } - } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportContext.java index 8b60e3d3dc..7cb923e9cc 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportContext.java @@ -25,6 +25,7 @@ import org.thingsboard.server.common.transport.TransportContext; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.coap.adaptors.JsonCoapAdaptor; import org.thingsboard.server.transport.coap.adaptors.ProtoCoapAdaptor; +import org.thingsboard.server.transport.coap.client.CoapClientContext; import org.thingsboard.server.transport.coap.efento.adaptor.EfentoCoapAdaptor; import java.util.concurrent.ConcurrentHashMap; @@ -52,6 +53,9 @@ public class CoapTransportContext extends TransportContext { @Autowired private EfentoCoapAdaptor efentoCoapAdaptor; + @Autowired + private CoapClientContext clientContext; + private final ConcurrentMap rpcAwaitingAck = new ConcurrentHashMap<>(); } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index d78f402e31..16d3aba7d6 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java @@ -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 tokenToCoapSessionInfoMap = new ConcurrentHashMap<>(); - private final ConcurrentMap sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); - private final Set rpcSubscriptions = ConcurrentHashMap.newKeySet(); - private final Set attributeSubscriptions = ConcurrentHashMap.newKeySet(); + private final ConcurrentMap sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); private final ConcurrentMap 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 coapObserveSessionInfos = sessionInfoToObserveRelationMap.keySet(); - Set 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 getCoapSessionInfoToObserveRelationMap() { - return sessionInfoToObserveRelationMap; - } - @Override protected void processHandleGet(CoapExchange exchange) { Optional 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 sessionToObserveRelationMap = CoapTransportResource.this.getCoapSessionInfoToObserveRelationMap(); - if (CoapTransportResource.this.getObserverCount() > 0 && !CollectionUtils.isEmpty(sessionToObserveRelationMap)) { - Optional 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); - } - } } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/OtaPackageTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/OtaPackageTransportResource.java index 74777d59d4..b629872924 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/OtaPackageTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/OtaPackageTransportResource.java @@ -24,9 +24,13 @@ import org.eclipse.californium.core.server.resources.CoapExchange; import org.eclipse.californium.core.server.resources.Resource; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.security.DeviceTokenCredentials; import org.thingsboard.server.common.transport.TransportServiceCallback; +import org.thingsboard.server.common.transport.auth.SessionInfoCreator; +import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.coap.callback.CoapDeviceAuthCallback; @@ -68,19 +72,21 @@ public class OtaPackageTransportResource extends AbstractCoapTransportResource { return; } transportService.process(DeviceTransportType.COAP, TransportProtos.ValidateDeviceTokenRequestMsg.newBuilder().setToken(credentials.get().getCredentialsId()).build(), - new CoapDeviceAuthCallback(transportContext, exchange, (sessionInfo, deviceProfile) -> { - getOtaPackageCallback(sessionInfo, exchange, otaPackageType); + new CoapDeviceAuthCallback(exchange, (msg, deviceProfile) -> { + getOtaPackageCallback(msg, exchange, otaPackageType); })); } - private void getOtaPackageCallback(TransportProtos.SessionInfoProto sessionInfo, CoapExchange exchange, OtaPackageType firmwareType) { + private void getOtaPackageCallback(ValidateDeviceCredentialsResponse msg, CoapExchange exchange, OtaPackageType firmwareType) { + TenantId tenantId = msg.getDeviceInfo().getTenantId(); + DeviceId deviceId = msg.getDeviceInfo().getDeviceId(); TransportProtos.GetOtaPackageRequestMsg requestMsg = TransportProtos.GetOtaPackageRequestMsg.newBuilder() - .setTenantIdMSB(sessionInfo.getTenantIdMSB()) - .setTenantIdLSB(sessionInfo.getTenantIdLSB()) - .setDeviceIdMSB(sessionInfo.getDeviceIdMSB()) - .setDeviceIdLSB(sessionInfo.getDeviceIdLSB()) + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) .setType(firmwareType.name()).build(); - transportContext.getTransportService().process(sessionInfo, requestMsg, new OtaPackageCallback(exchange)); + transportContext.getTransportService().process(SessionInfoCreator.create(msg, transportContext, UUID.randomUUID()), requestMsg, new OtaPackageCallback(exchange)); } private Optional decodeCredentials(Request request) { diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/TransportConfigurationContainer.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/TransportConfigurationContainer.java new file mode 100644 index 0000000000..39b8b2a4ca --- /dev/null +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/TransportConfigurationContainer.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.coap; + +import com.google.protobuf.Descriptors; +import com.google.protobuf.DynamicMessage; +import lombok.Data; + +@Data +public 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; + } +} diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java index 88b0be4634..54680df2bd 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java @@ -107,7 +107,7 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { Response response = new Response(CoAP.ResponseCode.CONTENT); JsonElement result = JsonConverter.toJson(msg); response.setPayload(result.toString()); - response.setAcknowledged(isConfirmable); + response.setConfirmable(isConfirmable); return response; } @@ -125,8 +125,8 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { public Response convertToPublish(boolean isConfirmable, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException { if (msg.getSharedStateMsg()) { if (StringUtils.isEmpty(msg.getError())) { - Response response = new Response(CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); - response.setAcknowledged(isConfirmable); + Response response = new Response(CoAP.ResponseCode.CONTENT); + response.setConfirmable(isConfirmable); TransportProtos.AttributeUpdateNotificationMsg notificationMsg = TransportProtos.AttributeUpdateNotificationMsg.newBuilder().addAllSharedUpdated(msg.getSharedAttributeListList()).build(); JsonObject result = JsonConverter.toJson(notificationMsg); response.setPayload(result.toString()); @@ -139,7 +139,7 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { return new Response(CoAP.ResponseCode.NOT_FOUND); } else { Response response = new Response(CoAP.ResponseCode.CONTENT); - response.setAcknowledged(isConfirmable); + response.setConfirmable(isConfirmable); JsonObject result = JsonConverter.toJson(msg); response.setPayload(result.toString()); return response; @@ -148,7 +148,7 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { } private Response getObserveNotification(boolean confirmable, JsonElement json) { - Response response = new Response(CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + Response response = new Response(CoAP.ResponseCode.CONTENT); response.setPayload(json.toString()); response.setConfirmable(confirmable); return response; diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java index 3f3a67ff3c..93d9a35029 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java @@ -130,8 +130,8 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor { public Response convertToPublish(boolean isConfirmable, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException { if (msg.getSharedStateMsg()) { if (StringUtils.isEmpty(msg.getError())) { - Response response = new Response(CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); - response.setAcknowledged(isConfirmable); + Response response = new Response(CoAP.ResponseCode.CONTENT); + response.setConfirmable(isConfirmable); TransportProtos.AttributeUpdateNotificationMsg notificationMsg = TransportProtos.AttributeUpdateNotificationMsg.newBuilder().addAllSharedUpdated(msg.getSharedAttributeListList()).build(); response.setPayload(notificationMsg.toByteArray()); return response; @@ -143,7 +143,7 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor { return new Response(CoAP.ResponseCode.NOT_FOUND); } else { Response response = new Response(CoAP.ResponseCode.CONTENT); - response.setAcknowledged(isConfirmable); + response.setConfirmable(isConfirmable); response.setPayload(msg.toByteArray()); return response; } @@ -151,9 +151,9 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor { } private Response getObserveNotification(boolean confirmable, byte[] notification) { - Response response = new Response(CoAP.ResponseCode._UNKNOWN_SUCCESS_CODE); + Response response = new Response(CoAP.ResponseCode.CONTENT); response.setPayload(notification); - response.setAcknowledged(confirmable); + response.setConfirmable(confirmable); return response; } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/AbstractSyncSessionCallback.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/AbstractSyncSessionCallback.java new file mode 100644 index 0000000000..b755b7097e --- /dev/null +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/AbstractSyncSessionCallback.java @@ -0,0 +1,74 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.coap.callback; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.californium.core.coap.Request; +import org.eclipse.californium.core.server.resources.CoapExchange; +import org.thingsboard.server.common.transport.SessionMsgListener; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.coap.client.TbCoapClientState; +import org.thingsboard.server.transport.coap.client.TbCoapObservationState; + +import java.util.UUID; + +@RequiredArgsConstructor +@Slf4j +public abstract class AbstractSyncSessionCallback implements SessionMsgListener { + + protected final TbCoapClientState state; + protected final CoapExchange exchange; + protected final Request request; + + @Override + public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg getAttributesResponse) { + logUnsupportedCommandMessage(getAttributesResponse); + } + + @Override + public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotification) { + logUnsupportedCommandMessage(attributeUpdateNotification); + } + + @Override + public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { + + } + + @Override + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest) { + logUnsupportedCommandMessage(toDeviceRequest); + } + + @Override + public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { + logUnsupportedCommandMessage(toServerResponse); + } + + private void logUnsupportedCommandMessage(Object update) { + log.trace("[{}] Ignore unsupported update: {}", state.getDeviceId(), update); + } + + public static boolean isConRequest(TbCoapObservationState state) { + if (state != null) { + return state.getExchange().advanced().getRequest().isConfirmable(); + } else { + return false; + } + } + +} diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapDeviceAuthCallback.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapDeviceAuthCallback.java index fab0c9b615..18fb5da216 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapDeviceAuthCallback.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapDeviceAuthCallback.java @@ -31,12 +31,10 @@ import java.util.function.BiConsumer; @Slf4j public class CoapDeviceAuthCallback implements TransportServiceCallback { - private final TransportContext transportContext; private final CoapExchange exchange; - private final BiConsumer onSuccess; + private final BiConsumer onSuccess; - public CoapDeviceAuthCallback(TransportContext transportContext, CoapExchange exchange, BiConsumer onSuccess) { - this.transportContext = transportContext; + public CoapDeviceAuthCallback(CoapExchange exchange, BiConsumer onSuccess) { this.exchange = exchange; this.onSuccess = onSuccess; } @@ -45,8 +43,7 @@ public class CoapDeviceAuthCallback implements TransportServiceCallback clients = new ConcurrentHashMap<>(); + private final ConcurrentMap clientsByToken = new ConcurrentHashMap<>(); + + @Override + public boolean registerAttributeObservation(TbCoapClientState clientState, String token, CoapExchange exchange) { + return registerFeatureObservation(clientState, token, exchange, FeatureType.ATTRIBUTES); + } + + @Override + public boolean registerRpcObservation(TbCoapClientState clientState, String token, CoapExchange exchange) { + return registerFeatureObservation(clientState, token, exchange, FeatureType.RPC); + } + + @Override + public void onUplink(TransportProtos.SessionInfoProto sessionInfo) { + getClientState(toDeviceId(sessionInfo)).updateLastUplinkTime(); + } + + @Override + public AtomicInteger getNotificationCounterByToken(String token) { + TbCoapClientState state = clientsByToken.get(token); + if (state == null) { + log.trace("Failed to find state using token: {}", token); + return null; + } + if (state.getAttrs() != null && state.getAttrs().getToken().equals(token)) { + return state.getAttrs().getObserveCounter(); + } else { + log.trace("Failed to find attr subscription using token: {}", token); + } + if (state.getRpc() != null && state.getRpc().getToken().equals(token)) { + return state.getRpc().getObserveCounter(); + } else { + log.trace("Failed to find rpc subscription using token: {}", token); + } + return null; + } + + @Override + public void registerObserveRelation(String token, ObserveRelation relation) { + TbCoapClientState state = clientsByToken.get(token); + if (state == null) { + log.trace("Failed to find state using token: {}", token); + return; + } + if (state.getAttrs() != null && state.getAttrs().getToken().equals(token)) { + state.getAttrs().setObserveRelation(relation); + } else { + log.trace("Failed to find attr subscription using token: {}", token); + } + if (state.getRpc() != null && state.getRpc().getToken().equals(token)) { + state.getRpc().setObserveRelation(relation); + } else { + log.trace("Failed to find rpc subscription using token: {}", token); + } + } + + @Override + public void deregisterObserveRelation(String token) { + TbCoapClientState state = clientsByToken.remove(token); + if (state == null) { + log.trace("Failed to find state using token: {}", token); + return; + } + if (state.getAttrs() != null && state.getAttrs().getToken().equals(token)) { + cancelAttributeSubscription(state); + } else { + log.trace("Failed to find attr subscription using token: {}", token); + } + if (state.getRpc() != null && state.getRpc().getToken().equals(token)) { + cancelRpcSubscription(state); + } else { + log.trace("Failed to find rpc subscription using token: {}", token); + } + } + + @Override + public void reportActivity() { + for (TbCoapClientState state : clients.values()) { + if (state.getSession() != null) { + transportService.reportActivity(state.getSession()); + } + } + } + + private boolean registerFeatureObservation(TbCoapClientState state, String token, CoapExchange exchange, FeatureType featureType) { + state.lock(); + try { + boolean newObservation; + if (FeatureType.ATTRIBUTES.equals(featureType)) { + if (state.getAttrs() == null) { + newObservation = true; + state.setAttrs(new TbCoapObservationState(exchange, token)); + } else { + newObservation = !state.getAttrs().getToken().equals(token); + if (newObservation) { + TbCoapObservationState old = state.getAttrs(); + state.setAttrs(new TbCoapObservationState(exchange, token)); + old.getExchange().respond(CoAP.ResponseCode.DELETED); + } + } + } else { + if (state.getRpc() == null) { + newObservation = true; + state.setRpc(new TbCoapObservationState(exchange, token)); + } else { + newObservation = !state.getRpc().getToken().equals(token); + if (newObservation) { + TbCoapObservationState old = state.getRpc(); + state.setRpc(new TbCoapObservationState(exchange, token)); + old.getExchange().respond(CoAP.ResponseCode.DELETED); + } + } + } + if (newObservation) { + clientsByToken.put(token, state); + if (state.getSession() == null) { + TransportProtos.SessionInfoProto session = SessionInfoCreator.create(state.getCredentials(), transportContext, UUID.randomUUID()); + state.setSession(session); + transportService.registerAsyncSession(session, new CoapSessionListener(state)); + transportService.process(session, getSessionEventMsg(TransportProtos.SessionEvent.OPEN), null); + } + if (FeatureType.ATTRIBUTES.equals(featureType)) { + transportService.process(state.getSession(), + TransportProtos.SubscribeToAttributeUpdatesMsg.getDefaultInstance(), new CoapNoOpCallback(exchange)); + transportService.process(state.getSession(), + TransportProtos.GetAttributeRequestMsg.newBuilder().setOnlyShared(true).build(), + new CoapNoOpCallback(exchange)); + } else { + transportService.process(state.getSession(), + TransportProtos.SubscribeToRPCMsg.getDefaultInstance(), + new CoapOkCallback(exchange, CoAP.ResponseCode.VALID, CoAP.ResponseCode.INTERNAL_SERVER_ERROR) + ); + } + } + return newObservation; + } finally { + state.unlock(); + } + } + + @Override + public void deregisterAttributeObservation(TbCoapClientState state, String token, CoapExchange exchange) { + state.lock(); + try { + clientsByToken.remove(token); + if (state.getSession() == null) { + log.trace("[{}] Failed to delete attribute observation: {}. Session is not present.", state.getDeviceId(), token); + return; + } + if (state.getAttrs() == null) { + log.trace("[{}] Failed to delete attribute observation: {}. It is not registered.", state.getDeviceId(), token); + return; + } + if (!state.getAttrs().getToken().equals(token)) { + log.trace("[{}] Failed to delete attribute observation: {}. Token mismatch.", state.getDeviceId(), token); + return; + } + cancelAttributeSubscription(state); + } finally { + state.unlock(); + } + } + + @Override + public void deregisterRpcObservation(TbCoapClientState state, String token, CoapExchange exchange) { + state.lock(); + try { + clientsByToken.remove(token); + if (state.getSession() == null) { + log.trace("[{}] Failed to delete rpc observation: {}. Session is not present.", state.getDeviceId(), token); + return; + } + if (state.getRpc() == null) { + log.trace("[{}] Failed to delete rpc observation: {}. It is not registered.", state.getDeviceId(), token); + return; + } + if (!state.getRpc().getToken().equals(token)) { + log.trace("[{}] Failed to delete rpc observation: {}. Token mismatch.", state.getDeviceId(), token); + return; + } + cancelRpcSubscription(state); + } finally { + state.unlock(); + } + } + + @Override + public TbCoapClientState getOrCreateClient(SessionMsgType type, ValidateDeviceCredentialsResponse deviceCredentials, DeviceProfile deviceProfile) throws AdaptorException { + DeviceId deviceId = deviceCredentials.getDeviceInfo().getDeviceId(); + TbCoapClientState state = getClientState(deviceId); + state.lock(); + try { + if (state.getConfiguration() == null || state.getAdaptor() == null) { + state.setConfiguration(getTransportConfigurationContainer(deviceProfile)); + state.setAdaptor(getCoapTransportAdaptor(state.getConfiguration().isJsonPayload())); + } + if (state.getCredentials() == null) { + state.setCredentials(deviceCredentials); + } + } finally { + state.unlock(); + } + return state; + } + + @Override + public TransportProtos.SessionInfoProto getNewSyncSession(TbCoapClientState state) { + return SessionInfoCreator.create(state.getCredentials(), transportContext, UUID.randomUUID()); + } + + private TbCoapClientState getClientState(DeviceId deviceId) { + return clients.computeIfAbsent(deviceId, TbCoapClientState::new); + } + + private static DeviceId toDeviceId(TransportProtos.SessionInfoProto s) { + return new DeviceId(new UUID(s.getDeviceIdMSB(), s.getDeviceIdLSB())); + } + + private static TransportProtos.SessionEventMsg getSessionEventMsg(TransportProtos.SessionEvent event) { + return TransportProtos.SessionEventMsg.newBuilder() + .setSessionType(TransportProtos.SessionType.ASYNC) + .setEvent(event).build(); + } + + 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(); + } + + @RequiredArgsConstructor + private class CoapSessionListener implements SessionMsgListener { + + private final TbCoapClientState state; + + @Override + public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg msg) { + TbCoapObservationState attrs = state.getAttrs(); + if (attrs != null) { + try { + boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getAttrs()); + Response response = state.getAdaptor().convertToPublish(conRequest, msg); + attrs.getExchange().respond(response); + } catch (AdaptorException e) { + log.trace("Failed to reply due to error", e); + cancelObserveRelation(attrs); + cancelAttributeSubscription(state); + } + } else { + log.debug("[{}] Get Attrs exchange is empty", state.getDeviceId()); + } + } + + @Override + public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg msg) { + log.trace("[{}] Received attributes update notification to device", sessionId); + TbCoapObservationState attrs = state.getAttrs(); + if (attrs != null) { + try { + boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getAttrs()); + Response response = state.getAdaptor().convertToPublish(conRequest, msg); + attrs.getExchange().respond(response); + } catch (AdaptorException e) { + log.trace("[{}] Failed to reply due to error", state.getDeviceId(), e); + cancelObserveRelation(attrs); + cancelAttributeSubscription(state); + } + } else { + log.debug("[{}] Get Attrs exchange is empty", state.getDeviceId()); + } + } + + @Override + public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { + log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); + cancelRpcSubscription(state); + cancelAttributeSubscription(state); + } + + @Override + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg msg) { + log.trace("[{}] Received RPC command to device", sessionId); + boolean sent = false; + boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getRpc()); + try { + Response response = state.getAdaptor().convertToPublish(conRequest, msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder()); + int requestId = getNextMsgId(); + response.setMID(requestId); + if (msg.getPersisted() && conRequest) { + transportContext.getRpcAwaitingAck().put(requestId, msg); + transportContext.getScheduler().schedule(() -> { + TransportProtos.ToDeviceRpcRequestMsg awaitingAckMsg = transportContext.getRpcAwaitingAck().remove(requestId); + if (awaitingAckMsg != null) { + transportService.process(state.getSession(), 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(state.getSession(), rpcRequestMsg, false, TransportServiceCallback.EMPTY); + } + })); + } + state.getRpc().getExchange().respond(response); + sent = true; + } catch (AdaptorException e) { + log.trace("Failed to reply due to error", e); + cancelObserveRelation(state.getRpc()); + cancelRpcSubscription(state); + } finally { + if (msg.getPersisted() && !conRequest) { + transportService.process(state.getSession(), msg, sent, TransportServiceCallback.EMPTY); + } + } + } + + @Override + public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg msg) { + + } + + private void cancelObserveRelation(TbCoapObservationState attrs) { + if (attrs.getObserveRelation() != null) { + attrs.getObserveRelation().cancel(); + } + } + } + + protected int getNextMsgId() { + return ThreadLocalRandom.current().nextInt(NONE, MAX_MID + 1); + } + + private void cancelRpcSubscription(TbCoapClientState state) { + if (state.getRpc() != null) { + clientsByToken.remove(state.getRpc().getToken()); + CoapExchange exchange = state.getAttrs().getExchange(); + state.setRpc(null); + transportService.process(state.getSession(), + TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), + new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); + if (state.getAttrs() == null) { + transportService.process(state.getSession(), getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); + transportService.deregisterSession(state.getSession()); + state.setSession(null); + } + } + } + + private void cancelAttributeSubscription(TbCoapClientState state) { + if (state.getAttrs() != null) { + clientsByToken.remove(state.getAttrs().getToken()); + CoapExchange exchange = state.getAttrs().getExchange(); + state.setAttrs(null); + transportService.process(state.getSession(), + TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setUnsubscribe(true).build(), + new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); + if (state.getRpc() == null) { + transportService.process(state.getSession(), getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); + transportService.deregisterSession(state.getSession()); + state.setSession(null); + } + } + } +} diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java new file mode 100644 index 0000000000..a2658ed82b --- /dev/null +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java @@ -0,0 +1,83 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.coap.client; + +import lombok.Data; +import lombok.Getter; +import lombok.Setter; +import org.eclipse.californium.core.network.Exchange; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.coap.TransportConfigurationContainer; +import org.thingsboard.server.transport.coap.adaptors.CoapTransportAdaptor; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +@Data +public class TbCoapClientState { + + private final DeviceId deviceId; + private final Lock lock; + + private volatile TransportConfigurationContainer configuration; + private volatile CoapTransportAdaptor adaptor; + private volatile ValidateDeviceCredentialsResponse credentials; + + private volatile TransportProtos.SessionInfoProto session; + + private volatile TbCoapObservationState attrs; + private volatile TbCoapObservationState rpc; + + @Getter + @Setter + private boolean asleep; + @Getter + private long lastUplinkTime; + @Getter + @Setter + private Future sleepTask; + + private boolean firstEdrxDownlink = true; + + public TbCoapClientState(DeviceId deviceId) { + this.deviceId = deviceId; + this.lock = new ReentrantLock(); + } + + public void lock() { + lock.lock(); + } + + public void unlock() { + lock.unlock(); + } + + public long updateLastUplinkTime(){ + this.lastUplinkTime = System.currentTimeMillis(); + this.firstEdrxDownlink = true; + return lastUplinkTime; + } + +} diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapObservationState.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapObservationState.java new file mode 100644 index 0000000000..5bfd71aed4 --- /dev/null +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapObservationState.java @@ -0,0 +1,34 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.coap.client; + +import lombok.Data; +import lombok.RequiredArgsConstructor; +import org.eclipse.californium.core.observe.ObserveRelation; +import org.eclipse.californium.core.server.resources.CoapExchange; + +import java.util.concurrent.atomic.AtomicInteger; + +@Data +@RequiredArgsConstructor +public class TbCoapObservationState { + + private final CoapExchange exchange; + private final String token; + private final AtomicInteger observeCounter = new AtomicInteger(0); + private volatile ObserveRelation observeRelation; + +} diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java index fe86d8794c..9cb1146cf2 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/efento/CoapEfentoTransportResource.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransp import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.EfentoCoapDeviceTypeConfiguration; import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.common.transport.auth.SessionInfoCreator; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.coap.MeasurementTypeProtos; import org.thingsboard.server.gen.transport.coap.MeasurementsProtos; @@ -81,7 +82,8 @@ public class CoapEfentoTransportResource extends AbstractCoapTransportResource { log.trace("Successfully parsed Efento ProtoMeasurements: [{}]", protoMeasurements.getCloudToken()); String token = protoMeasurements.getCloudToken(); transportService.process(DeviceTransportType.COAP, TransportProtos.ValidateDeviceTokenRequestMsg.newBuilder().setToken(token).build(), - new CoapDeviceAuthCallback(transportContext, exchange, (sessionInfo, deviceProfile) -> { + new CoapDeviceAuthCallback(exchange, (msg, deviceProfile) -> { + TransportProtos.SessionInfoProto sessionInfo = SessionInfoCreator.create(msg, transportContext, UUID.randomUUID()); UUID sessionId = new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()); try { validateEfentoTransportConfiguration(deviceProfile); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/AbstractLwM2mTransportResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/AbstractLwM2mTransportResource.java index 2c82facd23..f23b84488b 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/AbstractLwM2mTransportResource.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/AbstractLwM2mTransportResource.java @@ -43,50 +43,4 @@ public abstract class AbstractLwM2mTransportResource extends LwM2mCoapResource { protected abstract void processHandlePost(CoapExchange exchange); - public static class CoapOkCallback implements TransportServiceCallback { - private final CoapExchange exchange; - private final CoAP.ResponseCode onSuccessResponse; - private final CoAP.ResponseCode onFailureResponse; - - public CoapOkCallback(CoapExchange exchange, CoAP.ResponseCode onSuccessResponse, CoAP.ResponseCode onFailureResponse) { - this.exchange = exchange; - this.onSuccessResponse = onSuccessResponse; - this.onFailureResponse = onFailureResponse; - } - - @Override - public void onSuccess(Void msg) { - Response response = new Response(onSuccessResponse); - response.setAcknowledged(isConRequest()); - exchange.respond(response); - } - - @Override - public void onError(Throwable e) { - exchange.respond(onFailureResponse); - } - - private boolean isConRequest() { - return exchange.advanced().getRequest().isConfirmable(); - } - } - - public static class CoapNoOpCallback implements TransportServiceCallback { - private final CoapExchange exchange; - - CoapNoOpCallback(CoapExchange exchange) { - this.exchange = exchange; - } - - @Override - public void onSuccess(Void msg) { - } - - @Override - public void onError(Throwable e) { - exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); - } - } - - } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java index 7a70714b5a..070cff7666 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java @@ -161,6 +161,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl //TODO: check that the client supports FW and SW by checking the supported objects in the model. List attributesToFetch = new ArrayList<>(); LwM2MClientFwOtaInfo fwInfo = getOrInitFwInfo(client); + if (fwInfo.isSupported()) { attributesToFetch.add(FIRMWARE_TITLE); attributesToFetch.add(FIRMWARE_VERSION);