Browse Source

Uplink notifications for PSM & eDRX for CoAP in MSA deployment

pull/4966/head
Andrii Shvaika 5 years ago
parent
commit
99b19034e2
  1. 6
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  2. 11
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  3. 2
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java
  4. 3
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java
  5. 7
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  6. 3
      common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
  7. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java
  8. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java
  9. 3
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java
  10. 1
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

6
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -36,6 +36,12 @@ public class DataConstants {
public static final String ALARM_CONDITION_REPEATS = "alarmConditionRepeats";
public static final String ALARM_CONDITION_DURATION = "alarmConditionDuration";
public static final String PERSISTENT = "persistent";
public static final String COAP_TRANSPORT_NAME = "COAP";
public static final String LWM2M_TRANSPORT_NAME = "LWM2M";
public static final String MQTT_TRANSPORT_NAME = "MQTT";
public static final String HTTP_TRANSPORT_NAME = "HTTP";
public static final String SNMP_TRANSPORT_NAME = "SNMP";
public static final String[] allScopes() {
return new String[]{CLIENT_SCOPE, SHARED_SCOPE, SERVER_SCOPE};

11
common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

@ -73,6 +73,7 @@ public class HashPartitionService implements PartitionService {
private Map<String, TopicPartitionInfo> tbCoreNotificationTopics = new HashMap<>();
private Map<String, TopicPartitionInfo> tbRuleEngineNotificationTopics = new HashMap<>();
private Map<String, List<ServiceInfo>> tbTransportServicesByType = new HashMap<>();
private List<ServiceInfo> currentOtherServices;
private HashFunction hashFunction;
@ -127,6 +128,7 @@ public class HashPartitionService implements PartitionService {
@Override
public synchronized void recalculatePartitions(ServiceInfo currentService, List<ServiceInfo> otherServices) {
tbTransportServicesByType.clear();
logServiceInfo(currentService);
otherServices.forEach(this::logServiceInfo);
Map<ServiceQueueKey, List<ServiceInfo>> queueServicesMap = new HashMap<>();
@ -229,6 +231,12 @@ public class HashPartitionService implements PartitionService {
return Math.abs(hash % partitions);
}
@Override
public int countTransportsByType(String type) {
var list = tbTransportServicesByType.get(type);
return list == null ? 0 : list.size();
}
private Map<ServiceQueueKey, List<ServiceInfo>> getServiceKeyListMap(List<ServiceInfo> services) {
final Map<ServiceQueueKey, List<ServiceInfo>> currentMap = new HashMap<>();
services.forEach(serviceInfo -> {
@ -332,6 +340,9 @@ public class HashPartitionService implements PartitionService {
queueServiceList.computeIfAbsent(serviceQueueKey, key -> new ArrayList<>()).add(instance);
}
}
for (String transportType : instance.getTransportsList()) {
tbTransportServicesByType.computeIfAbsent(transportType, t -> new ArrayList<>()).add(instance);
}
}
private ServiceInfo resolveByPartitionIdx(List<ServiceInfo> servers, Integer partitionIdx) {

2
common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java

@ -59,4 +59,6 @@ public interface PartitionService {
TopicPartitionInfo getNotificationsTopic(ServiceType serviceType, String serviceId);
int resolvePartitionIndex(UUID entityId, int partitions);
int countTransportsByType(String type);
}

3
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java

@ -22,6 +22,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.coapserver.CoapServerService;
import org.thingsboard.server.coapserver.TbCoapServerComponent;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.transport.coap.efento.CoapEfentoTransportResource;
@ -72,6 +73,6 @@ public class CoapTransportService implements TbTransportService {
@Override
public String getName() {
return "COAP";
return DataConstants.COAP_TRANSPORT_NAME;
}
}

7
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java

@ -18,14 +18,13 @@ 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;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service;
import org.thingsboard.server.coapserver.CoapServerContext;
import org.thingsboard.server.coapserver.TbCoapServerComponent;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceTransportType;
@ -51,6 +50,7 @@ import org.thingsboard.server.common.transport.adaptor.AdaptorException;
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.queue.discovery.PartitionService;
import org.thingsboard.server.transport.coap.CoapTransportContext;
import org.thingsboard.server.transport.coap.TbCoapMessageObserver;
import org.thingsboard.server.transport.coap.TransportConfigurationContainer;
@ -81,6 +81,7 @@ public class DefaultCoapClientContext implements CoapClientContext {
private final CoapTransportContext transportContext;
private final TransportService transportService;
private final TransportDeviceProfileCache profileCache;
private final PartitionService partitionService;
private final ConcurrentMap<DeviceId, TbCoapClientState> clients = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TbCoapClientState> clientsByToken = new ConcurrentHashMap<>();
@ -214,7 +215,7 @@ public class DefaultCoapClientContext implements CoapClientContext {
return null;
}, timeout, TimeUnit.MILLISECONDS);
client.setSleepTask(task);
if (notifyOtherServers) {
if (notifyOtherServers && partitionService.countTransportsByType(DataConstants.COAP_TRANSPORT_NAME) > 1) {
transportService.notifyAboutUplink(getNewSyncSession(client), TransportProtos.UplinkNotificationMsg.newBuilder().setUplinkTs(uplinkTime).build(), TransportServiceCallback.EMPTY);
}
} finally {

3
common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java

@ -34,6 +34,7 @@ import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.data.id.DeviceId;
@ -436,7 +437,7 @@ public class DeviceApiController implements TbTransportService {
@Override
public String getName() {
return "HTTP";
return DataConstants.HTTP_TRANSPORT_NAME;
}
}

3
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java

@ -28,6 +28,7 @@ import org.eclipse.leshan.server.californium.registration.CaliforniumRegistratio
import org.eclipse.leshan.server.model.LwM2mModelProvider;
import org.springframework.stereotype.Component;
import org.thingsboard.server.cache.ota.OtaPackageDataCache;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MAuthorizer;
@ -177,7 +178,7 @@ public class DefaultLwM2mTransportService implements LwM2MTransportService {
@Override
public String getName() {
return "LWM2M";
return DataConstants.LWM2M_TRANSPORT_NAME;
}
}

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

@ -28,6 +28,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.TbTransportService;
import javax.annotation.PostConstruct;
@ -114,6 +115,6 @@ public class MqttTransportService implements TbTransportService {
@Override
public String getName() {
return "MQTT";
return DataConstants.MQTT_TRANSPORT_NAME;
}
}

3
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java

@ -35,6 +35,7 @@ import org.snmp4j.transport.DefaultUdpTransportMapping;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.data.kv.DataType;
import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec;
@ -300,7 +301,7 @@ public class SnmpTransportService implements TbTransportService {
@Override
public String getName() {
return "SNMP";
return DataConstants.SNMP_TRANSPORT_NAME;
}
@PreDestroy

1
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -573,7 +573,6 @@ public class DefaultTransportService implements TransportService {
@Override
public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback);

Loading…
Cancel
Save