Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard

pull/4963/head
Igor Kulikov 5 years ago
parent
commit
8468c33e14
  1. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
  2. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java
  3. 6
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java
  4. 5
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java
  5. 8
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/AbstractSyncSessionCallback.java
  6. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java
  7. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java
  8. 37
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  9. 7
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java
  10. 33
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapContentFormatUtil.java
  11. 14
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportCoapResource.java
  12. 1
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java
  13. 19
      pom.xml
  14. 2
      tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java
  15. 2
      tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java

2
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<TbCoapClientState, ObserveRelation> sessionInfoToObserveRelationMap = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TbCoapDtlsSessionInfo> dtlsSessionIdMap;
private final long timeout;
private final CoapClientContext clients;

2
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();
}

6
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;
}
}

5
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;
}
}

8
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);
}
}

2
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));

2
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);

37
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<DeviceProfile> 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);
}
}

7
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;

33
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;
}
}

14
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);
}

1
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());

19
pom.xml

@ -39,10 +39,11 @@
<javax-annotation.version>1.3.2</javax-annotation.version>
<jakarta.xml.bind-api.version>2.3.2</jakarta.xml.bind-api.version>
<jaxb-runtime.version>2.3.2</jaxb-runtime.version>
<spring-boot.version>2.3.9.RELEASE</spring-boot.version>
<spring.version>5.2.10.RELEASE</spring.version>
<spring-boot.version>2.3.12.RELEASE</spring-boot.version>
<spring.version>5.2.16.RELEASE</spring.version>
<spring-redis.version>5.2.11.RELEASE</spring-redis.version>
<spring-security.version>5.4.1</spring-security.version>
<spring-data-redis.version>2.4.1</spring-data-redis.version>
<spring-data-redis.version>2.4.3</spring-data-redis.version>
<jedis.version>3.3.0</jedis.version>
<jjwt.version>0.7.0</jjwt.version>
<json-path.version>2.2.0</json-path.version>
@ -62,7 +63,7 @@
<guava.version>28.2-jre</guava.version>
<caffeine.version>2.6.1</caffeine.version>
<commons-lang3.version>3.4</commons-lang3.version>
<commons-io.version>2.5</commons-io.version>
<commons-io.version>2.11.0</commons-io.version>
<commons-csv.version>1.4</commons-csv.version>
<jackson.version>2.12.1</jackson.version>
<jackson-annotations.version>2.12.1</jackson-annotations.version>
@ -79,7 +80,7 @@
<grpc.version>1.38.0</grpc.version>
<lombok.version>1.18.18</lombok.version>
<paho.client.version>1.2.4</paho.client.version>
<netty.version>4.1.60.Final</netty.version>
<netty.version>4.1.66.Final</netty.version>
<os-maven-plugin.version>1.7.0</os-maven-plugin.version>
<rabbitmq.version>4.8.0</rabbitmq.version>
<surfire.version>2.19.1</surfire.version>
@ -115,7 +116,7 @@
<micrometer.version>1.5.2</micrometer.version>
<protobuf-dynamic.version>1.0.3TB</protobuf-dynamic.version>
<wire-schema.version>3.4.0</wire-schema.version>
<twilio.version>7.54.2</twilio.version>
<twilio.version>8.17.0</twilio.version>
<hibernate-validator.version>6.0.13.Final</hibernate-validator.version>
<javax.el.version>3.0.0</javax.el.version>
<javax.validation-api.version>2.0.1.Final</javax.validation-api.version>
@ -1504,7 +1505,7 @@
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-redis</artifactId>
<version>${spring.version}</version>
<version>${spring-redis.version}</version>
</dependency>
<dependency>
<groupId>redis.clients</groupId>
@ -1678,6 +1679,10 @@
<groupId>com.github.spotbugs</groupId>
<artifactId>spotbugs-annotations</artifactId>
</exclusion>
<exclusion>
<groupId>commons-io</groupId>
<artifactId>commons-io</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>

2
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()) {

2
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()) {

Loading…
Cancel
Save