From 99b19034e202611bd3410bcfa5039f125154f19a Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 23 Jul 2021 13:50:48 +0300 Subject: [PATCH] Uplink notifications for PSM & eDRX for CoAP in MSA deployment --- .../thingsboard/server/common/data/DataConstants.java | 6 ++++++ .../server/queue/discovery/HashPartitionService.java | 11 +++++++++++ .../server/queue/discovery/PartitionService.java | 2 ++ .../server/transport/coap/CoapTransportService.java | 3 ++- .../coap/client/DefaultCoapClientContext.java | 7 ++++--- .../server/transport/http/DeviceApiController.java | 3 ++- .../lwm2m/server/DefaultLwM2mTransportService.java | 3 ++- .../server/transport/mqtt/MqttTransportService.java | 3 ++- .../transport/snmp/service/SnmpTransportService.java | 3 ++- .../transport/service/DefaultTransportService.java | 1 - 10 files changed, 33 insertions(+), 9 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 3b25bd1290..329e74d702 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/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}; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index c833b314e6..f00c13bcd5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -73,6 +73,7 @@ public class HashPartitionService implements PartitionService { private Map tbCoreNotificationTopics = new HashMap<>(); private Map tbRuleEngineNotificationTopics = new HashMap<>(); + private Map> tbTransportServicesByType = new HashMap<>(); private List currentOtherServices; private HashFunction hashFunction; @@ -127,6 +128,7 @@ public class HashPartitionService implements PartitionService { @Override public synchronized void recalculatePartitions(ServiceInfo currentService, List otherServices) { + tbTransportServicesByType.clear(); logServiceInfo(currentService); otherServices.forEach(this::logServiceInfo); Map> 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> getServiceKeyListMap(List services) { final Map> 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 servers, Integer partitionIdx) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java index 20c59378e8..be6752cfed 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java +++ b/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); } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java index 93096c901e..66fd870dc9 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java +++ b/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; } } 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 0182badfce..757ec45961 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,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 clients = new ConcurrentHashMap<>(); private final ConcurrentMap 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 { diff --git a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java index 874dd6bdd5..9c258ef00f 100644 --- a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java +++ b/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; } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java index 2d3d0c97da..356548b2bd 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java +++ b/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; } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index 5498752151..89709da405 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/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; } } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java index 3c296cc4ae..d0cee4a329 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java +++ b/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 diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 93f96795cc..eda114e498 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/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 callback) { - if (checkLimits(sessionInfo, msg, callback)) { reportActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback);