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 59aae12777..cf8c40a51c 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 @@ -67,8 +67,6 @@ 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 sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); - private final ConcurrentMap dtlsSessionIdMap; private final long timeout; private final CoapClientContext clients; diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java index 00374f6916..b38ac3c322 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java @@ -49,4 +49,6 @@ public interface CoapTransportAdaptor { ProvisionDeviceRequestMsg convertToProvisionRequestMsg(UUID sessionId, Request inbound) throws AdaptorException; + int getContentFormat(); + } 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 a1ef8ba205..e24e14faab 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 @@ -23,6 +23,7 @@ import com.google.protobuf.Descriptors; import com.google.protobuf.DynamicMessage; import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.core.coap.CoAP; +import org.eclipse.californium.core.coap.MediaTypeRegistry; import org.eclipse.californium.core.coap.Request; import org.eclipse.californium.core.coap.Response; import org.springframework.stereotype.Component; @@ -164,4 +165,9 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor { return payload; } + @Override + public int getContentFormat() { + return MediaTypeRegistry.APPLICATION_JSON; + } + } 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 93d9a35029..1c8f84fbd0 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 @@ -23,6 +23,7 @@ import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.util.JsonFormat; import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.core.coap.CoAP; +import org.eclipse.californium.core.coap.MediaTypeRegistry; import org.eclipse.californium.core.coap.Request; import org.eclipse.californium.core.coap.Response; import org.springframework.stereotype.Component; @@ -162,4 +163,8 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor { return JsonFormat.printer().includingDefaultValueFields().print(dynamicMessage); } + @Override + public int getContentFormat() { + return MediaTypeRegistry.APPLICATION_OCTET_STREAM; + } } 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 index b755b7097e..1eb6b28706 100644 --- 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 @@ -17,11 +17,14 @@ package org.thingsboard.server.transport.coap.callback; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.eclipse.californium.core.coap.MediaTypeRegistry; import org.eclipse.californium.core.coap.Request; +import org.eclipse.californium.core.coap.Response; 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.TbCoapContentFormatUtil; import org.thingsboard.server.transport.coap.client.TbCoapObservationState; import java.util.UUID; @@ -71,4 +74,9 @@ public abstract class AbstractSyncSessionCallback implements SessionMsgListener } } + protected void respond(Response response) { + response.getOptions().setContentFormat(TbCoapContentFormatUtil.getContentFormat(exchange.getRequestOptions().getContentFormat(), state.getContentFormat())); + exchange.respond(response); + } + } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java index 9c05098ce4..f6e5b52ff7 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java @@ -34,7 +34,7 @@ public class GetAttributesSyncSessionCallback extends AbstractSyncSessionCallbac @Override public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg msg) { try { - exchange.respond(state.getAdaptor().convertToPublish(AbstractSyncSessionCallback.isConRequest(state.getAttrs()), msg)); + respond(state.getAdaptor().convertToPublish(request.isConfirmable(), msg)); } catch (AdaptorException e) { log.trace("[{}] Failed to reply due to error", state.getDeviceId(), e); exchange.respond(new Response(CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java index b51758b51e..fb4b855f39 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java @@ -33,7 +33,7 @@ public class ToServerRpcSyncSessionCallback extends AbstractSyncSessionCallback @Override public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { try { - exchange.respond(state.getAdaptor().convertToPublish(isConRequest(state.getRpc()), toServerResponse)); + respond(state.getAdaptor().convertToPublish(request.isConfirmable(), toServerResponse)); } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index b5495042cf..b444a493d3 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -18,6 +18,7 @@ package org.thingsboard.server.transport.coap.client; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.core.coap.CoAP; +import org.eclipse.californium.core.coap.MediaTypeRegistry; import org.eclipse.californium.core.coap.Response; import org.eclipse.californium.core.observe.ObserveRelation; import org.eclipse.californium.core.server.resources.CoapExchange; @@ -328,8 +329,7 @@ public class DefaultCoapClientContext implements CoapClientContext { state.lock(); try { if (state.getConfiguration() == null || state.getAdaptor() == null) { - state.setConfiguration(getTransportConfigurationContainer(deviceProfile)); - state.setAdaptor(getCoapTransportAdaptor(state.getConfiguration().isJsonPayload())); + initStateAdaptor(deviceProfile, state); } if (state.getCredentials() == null) { state.init(deviceCredentials); @@ -393,6 +393,12 @@ public class DefaultCoapClientContext implements CoapClientContext { } } + private void initStateAdaptor(DeviceProfile deviceProfile, TbCoapClientState state) throws AdaptorException { + state.setConfiguration(getTransportConfigurationContainer(deviceProfile)); + state.setAdaptor(getCoapTransportAdaptor(state.getConfiguration().isJsonPayload())); + state.setContentFormat(state.getAdaptor().getContentFormat()); + } + private CoapTransportAdaptor getCoapTransportAdaptor(boolean jsonPayloadType) { return jsonPayloadType ? transportContext.getJsonCoapAdaptor() : transportContext.getProtoCoapAdaptor(); } @@ -409,7 +415,7 @@ public class DefaultCoapClientContext implements CoapClientContext { try { boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getAttrs()); Response response = state.getAdaptor().convertToPublish(conRequest, msg); - attrs.getExchange().respond(response); + respond(attrs.getExchange(), response, state.getContentFormat()); } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); cancelObserveRelation(attrs); @@ -440,10 +446,10 @@ public class DefaultCoapClientContext implements CoapClientContext { int requestId = getNextMsgId(); Response response = state.getAdaptor().convertToPublish(conRequest, msg); response.setMID(requestId); - attrs.getExchange().respond(response); if (conRequest) { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state))); } + respond(attrs.getExchange(), response, state.getContentFormat()); } catch (AdaptorException e) { log.trace("[{}] Failed to reply due to error", state.getDeviceId(), e); cancelObserveRelation(attrs); @@ -454,8 +460,24 @@ public class DefaultCoapClientContext implements CoapClientContext { } } + @Override + public void onDeviceProfileUpdate(TransportProtos.SessionInfoProto newSessionInfo, DeviceProfile deviceProfile) { + try { + initStateAdaptor(deviceProfile, state); + } catch (AdaptorException e) { + log.warn("[{}] Failed to update device profile: ", deviceProfile.getId(), e); + } + } + @Override public void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional deviceProfileOpt) { + if (deviceProfileOpt.isPresent()) { + try { + initStateAdaptor(deviceProfileOpt.get(), state); + } catch (AdaptorException e) { + log.warn("[{}] Failed to update device: ", device.getId(), e); + } + } state.onDeviceUpdate(device); } @@ -503,7 +525,7 @@ public class DefaultCoapClientContext implements CoapClientContext { if (conRequest) { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state))); } - state.getRpc().getExchange().respond(response); + respond(state.getRpc().getExchange(), response, state.getContentFormat()); sent = true; } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); @@ -705,4 +727,9 @@ public class DefaultCoapClientContext implements CoapClientContext { state.setAdaptor(null); //TODO: add optimistic lock check that the client was already deleted and cleanup "clients" map. } + + private void respond(CoapExchange exchange, Response response, int defContentFormat) { + response.getOptions().setContentFormat(TbCoapContentFormatUtil.getContentFormat(exchange.getRequestOptions().getContentFormat(), defContentFormat)); + exchange.respond(response); + } } 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 index 7613dc2ece..f106f57961 100644 --- 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 @@ -18,26 +18,21 @@ package org.thingsboard.server.transport.coap.client; import lombok.Data; import lombok.Getter; import lombok.Setter; -import org.eclipse.leshan.server.registration.Registration; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.device.data.CoapDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.data.PowerMode; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; -import org.thingsboard.server.common.data.id.TenantId; 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.ArrayList; import java.util.HashMap; import java.util.HashSet; -import java.util.List; import java.util.Map; import java.util.Set; -import java.util.UUID; import java.util.concurrent.Future; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -55,6 +50,8 @@ public class TbCoapClientState { private volatile DefaultCoapClientContext.CoapSessionListener listener; private volatile TbCoapObservationState attrs; private volatile TbCoapObservationState rpc; + private volatile int contentFormat; + private TransportProtos.AttributeUpdateNotificationMsg missedAttributeUpdates; private DeviceProfileId profileId; diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapContentFormatUtil.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapContentFormatUtil.java new file mode 100644 index 0000000000..29efb8f97b --- /dev/null +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapContentFormatUtil.java @@ -0,0 +1,33 @@ +/** + * 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 org.eclipse.californium.core.coap.MediaTypeRegistry; + +public class TbCoapContentFormatUtil { + + public static int getContentFormat(int requestFormat, int adaptorFormat) { + if (isStrict(adaptorFormat)) { + return adaptorFormat; + } else { + return requestFormat != MediaTypeRegistry.UNDEFINED ? requestFormat : adaptorFormat; + } + } + + public static boolean isStrict(int contentFormat) { + return contentFormat == MediaTypeRegistry.APPLICATION_OCTET_STREAM; + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportCoapResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportCoapResource.java index bbeda58270..bde1725807 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportCoapResource.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportCoapResource.java @@ -137,17 +137,17 @@ public class LwM2mTransportCoapResource extends AbstractLwM2mTransportResource { ); UUID currentId = UUID.fromString(idStr); Response response = new Response(CoAP.ResponseCode.CONTENT); - byte[] fwData = this.getOtaData(currentId); - log.debug("Read softWare data (length): [{}]", fwData.length); - if (fwData != null && fwData.length > 0) { - response.setPayload(fwData); + byte[] otaData = this.getOtaData(currentId); + log.debug("Read ota data (length): [{}]", otaData.length); + if (otaData.length > 0) { + response.setPayload(otaData); if (exchange.getRequestOptions().getBlock2() != null) { int chunkSize = exchange.getRequestOptions().getBlock2().getSzx(); - boolean lastFlag = fwData.length <= chunkSize; + boolean lastFlag = otaData.length <= chunkSize; response.getOptions().setBlock2(chunkSize, lastFlag, 0); - log.trace("With block2 Send currentId: [{}], length: [{}], chunkSize [{}], moreFlag [{}]", currentId.toString(), fwData.length, chunkSize, lastFlag); + log.trace("With block2 Send currentId: [{}], length: [{}], chunkSize [{}], moreFlag [{}]", currentId.toString(), otaData.length, chunkSize, lastFlag); } else { - log.trace("With block1 Send currentId: [{}], length: [{}], ", currentId.toString(), fwData.length); + log.trace("With block1 Send currentId: [{}], length: [{}], ", currentId.toString(), otaData.length); } exchange.respond(response); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java index ebf1021c41..c81d677155 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java @@ -317,7 +317,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl this.updateResourcesValue(lwM2MClient, lwM2mResource, path); } } - clientContext.update(lwM2MClient); if (clientContext.awake(lwM2MClient)) { // clientContext.awake calls clientContext.update log.debug("[{}] Device is awake", lwM2MClient.getEndpoint()); diff --git a/pom.xml b/pom.xml index 6d54fcdaa5..d0c9864ef3 100755 --- a/pom.xml +++ b/pom.xml @@ -39,10 +39,11 @@ 1.3.2 2.3.2 2.3.2 - 2.3.9.RELEASE - 5.2.10.RELEASE + 2.3.12.RELEASE + 5.2.16.RELEASE + 5.2.11.RELEASE 5.4.1 - 2.4.1 + 2.4.3 3.3.0 0.7.0 2.2.0 @@ -62,7 +63,7 @@ 28.2-jre 2.6.1 3.4 - 2.5 + 2.11.0 1.4 2.12.1 2.12.1 @@ -79,7 +80,7 @@ 1.38.0 1.18.18 1.2.4 - 4.1.60.Final + 4.1.66.Final 1.7.0 4.8.0 2.19.1 @@ -115,7 +116,7 @@ 1.5.2 1.0.3TB 3.4.0 - 7.54.2 + 8.17.0 6.0.13.Final 3.0.0 2.0.1.Final @@ -1504,7 +1505,7 @@ org.springframework.integration spring-integration-redis - ${spring.version} + ${spring-redis.version} redis.clients @@ -1678,6 +1679,10 @@ com.github.spotbugs spotbugs-annotations + + commons-io + commons-io + diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java index 53caeaeb0d..ca56f0e32e 100644 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java @@ -43,7 +43,7 @@ public class DictionaryParser { return line.startsWith("COPY public.ts_kv_dictionary ("); } - private void parseDictionaryDump(LineIterator iterator) { + private void parseDictionaryDump(LineIterator iterator) throws IOException { try { String tempLine; while (iterator.hasNext()) { diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java index 7c1b29f2af..173d4e8054 100644 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java @@ -58,7 +58,7 @@ public class RelatedEntitiesParser { return StringUtils.isBlank(line) || line.equals("\\."); } - private void processAllTables(LineIterator lineIterator) { + private void processAllTables(LineIterator lineIterator) throws IOException { String currentLine; try { while (lineIterator.hasNext()) {