?z@_Z+$B-|qW_Z)r9neq15Q+$S*@ouVGy6Uoqa-s0?0PGVs}CS^)b&A}QaiYJVC!C(R2abPjUMZuJ!q;`pkhnf^E
zN-qO{^xDO3@~c|YOMM}#z@X$6w1UVY2e$8IFB>!Ot{hz`l=Z
zp{@P$&-fzjc-2^w`0ZadoET#Kc~z!(jfOwMlmpD`_Npom$!1yfra^0cXl)IoroW|!
zk90?jdA80{`mRyU&@Go)fqSbfR0QZXdi)>Of~8C6ZrJP!&D*8-hQ~cGd8dCaeU8)@
zZxCZ#N~h$PvV@G(@gP8Cp^dpic2aeZ{Csq?C2BX|PSdCIcRcf26E$2f?X1q@|76)A
zOqb3U399Ngzb-qbg{M=c*t+INmM{=?9ebISEyZM`QW9ZUr!DPvz_xodb)8yhJaXp}
z@P?IylaD}NB~{%3pMJh9`i@6Xl|;DZCgdsQ@-LCRJg`2ug5xbn6DGPrRy-C
zP*ZKbOZ+(Bm^91@8~-TBEhhWKc#Dn6KeW!`ZC;fMLr>CYCH0nJ&1+}ITfR<_?>Q=#
zYIw%Dlqql|VZ77lerpB}yJWMI*LlVCGbo9~+}Sj%Qej#)mS)t=Sob}}^YF{9fpW5%
zof6)_fGquH1JO16ki^GQEsXI~FlecO-i@L^blcAM!q`{xMENiDI-lZCT4~o2DNr;E
z)Ad-}lU|byt0l(y7<#;uj^2!3!;0*4xb_0XWkSI?-xAIcD_ydn3Q@a)KN^WOD-(9
zJ?rIwzMR9O>>oD4yJTjGJ)o(nI*K%t%c@c%Nqr~^ke6{yCZca8dmd}}gnEo9++vtq
z_8nxn5TY_h|`X_#(f8WGzxxA5R#4(2lj3peMptJ(ihw60!a9W%wB
zrHX!YwsSBUT|Se2QunC^vq++%kzUQd(KYb8&tK>q@K+L@mRQsbI*pYUis|U%H768F
zY9Z5`Nn>YY`*3k6hdRyrjwC+j=4eZ$ckle|;$)wugtM-_#>z2Yi
z^WVE=&ZHv-$%%j`h(q9vWqmEz7s>dr{duamIU5_~Wp@N&n=o$h;0;WY%I4iPvea*iZRphgBTL22a?4qdaVs7GMaoxd#{{d^`RM-Fj
literal 0
HcmV?d00001
diff --git a/ui-ngx/src/assets/white_label_logo.svg b/ui-ngx/src/assets/white_label_logo.svg
new file mode 100644
index 0000000000..706ca443f0
--- /dev/null
+++ b/ui-ngx/src/assets/white_label_logo.svg
@@ -0,0 +1 @@
+
\ No newline at end of file
From 8f2438d6ab57ca22b940a001e9a05e68225e6c69 Mon Sep 17 00:00:00 2001
From: Viacheslav Klimov
Date: Mon, 22 Mar 2021 17:17:42 +0200
Subject: [PATCH 004/476] SNMP devices balancing (#4254)
* Fix merge errors
* Implement SNMP transports balancing
* Refactor; implement transport device cache
* Refactor
* Finish up device lifecycle handling implementing; refactor
* Refactor
* Change base image to thingsboard/openjdk11 for msa snmp transport
* Refactor
* Change transport services names to upper-case
---
.../actors/service/DefaultActorService.java | 4 +-
.../DefaultTbApiUsageStateService.java | 2 +-
.../apiusage/TbApiUsageStateService.java | 2 +-
.../queue/DefaultTbCoreConsumerService.java | 2 +-
.../DefaultTbRuleEngineConsumerService.java | 2 +-
.../service/queue/TbCoreConsumerService.java | 2 +-
.../queue/TbRuleEngineConsumerService.java | 2 +-
.../processing/AbstractConsumerService.java | 2 +-
.../state/DefaultDeviceStateService.java | 2 +-
.../service/state/DeviceStateService.java | 2 +-
.../DefaultSubscriptionManagerService.java | 2 +-
.../DefaultTbLocalSubscriptionService.java | 5 +-
.../SubscriptionManagerService.java | 2 +-
.../TbLocalSubscriptionService.java | 4 +-
.../AbstractSubscriptionService.java | 19 +-
.../telemetry/AlarmSubscriptionService.java | 3 +-
.../DefaultAlarmSubscriptionService.java | 23 -
.../TelemetrySubscriptionService.java | 3 +-
.../transport/DefaultTransportApiService.java | 94 ++++-
.../server/dao/device/DeviceService.java | 3 +
.../common/data/TbTransportService.java | 20 +
.../data/DeviceTransportConfiguration.java | 4 +-
.../SnmpProfileTransportConfiguration.java | 10 +-
.../DefaultTbServiceInfoProvider.java | 20 +
.../queue/discovery/HashPartitionService.java | 20 +-
.../queue/discovery/PartitionService.java | 4 +
.../discovery/TbApplicationEventListener.java | 1 +
.../queue/discovery/ZkDiscoveryService.java | 12 +-
.../ClusterTopologyChangeEvent.java | 3 +-
.../{ => event}/PartitionChangeEvent.java | 3 +-
.../event/ServiceListChangedEvent.java | 35 ++
.../{ => event}/TbApplicationEvent.java | 2 +-
.../queue/util/TbSnmpTransportComponent.java | 29 ++
common/queue/src/main/proto/queue.proto | 35 ++
.../transport/coap/CoapTransportService.java | 9 +-
.../transport/http/DeviceApiController.java | 10 +-
.../lwm2m/server/LwM2mTransportService.java | 3 +-
.../server/LwM2mTransportServiceImpl.java | 6 +
.../transport/mqtt/MqttTransportService.java | 8 +-
common/transport/snmp/pom.xml | 4 -
.../transport/snmp/SnmpDeviceSimulator.java | 81 ++++
.../transport/snmp/SnmpTransportContext.java | 311 +++++++++++---
.../transport/snmp/SnmpTransportService.java | 147 -------
.../ServiceListChangedEventListener.java | 35 ++
.../event/SnmpTransportListChangedEvent.java | 24 ++
...SnmpTransportListChangedEventListener.java | 34 ++
.../service/ProtoTransportEntityService.java | 92 ++++
.../SnmpTransportBalancingService.java | 92 ++++
.../SnmpTransportService.java} | 149 +++++--
.../snmp/session/DeviceSessionContext.java | 166 ++++++++
.../snmp/session/DeviceSessionCtx.java | 206 ---------
.../common/transport/DeviceUpdatedEvent.java | 28 ++
.../common/transport/SessionMsgListener.java | 4 +
.../common/transport/TransportContext.java | 1 +
.../common/transport/TransportService.java | 12 +
.../DefaultTransportDeviceProfileCache.java | 4 +-
.../service/DefaultTransportService.java | 77 +++-
.../server/dao/device/DeviceDao.java | 4 +
.../server/dao/device/DeviceServiceImpl.java | 7 +
.../dao/sql/device/DeviceRepository.java | 6 +
.../server/dao/sql/device/JpaDeviceDao.java | 6 +
docker/tb-transports/snmp/conf/logback.xml | 50 +++
.../snmp/conf/tb-snmp-transport.conf | 23 +
msa/transport/snmp/docker/Dockerfile | 2 +-
transport/snmp/src/main/conf/logback.xml | 45 ++
.../snmp/src/main/conf/tb-snmp-transport.conf | 22 +
transport/snmp/src/main/resources/logback.xml | 36 ++
.../src/main/resources/tb-snmp-transport.yml | 394 ++++++++++++++++++
68 files changed, 1923 insertions(+), 553 deletions(-)
create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/TbTransportService.java
rename common/queue/src/main/java/org/thingsboard/server/queue/discovery/{ => event}/ClusterTopologyChangeEvent.java (91%)
rename common/queue/src/main/java/org/thingsboard/server/queue/discovery/{ => event}/PartitionChangeEvent.java (93%)
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ServiceListChangedEvent.java
rename common/queue/src/main/java/org/thingsboard/server/queue/discovery/{ => event}/TbApplicationEvent.java (95%)
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/util/TbSnmpTransportComponent.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulator.java
delete mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportService.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/event/ServiceListChangedEventListener.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/event/SnmpTransportListChangedEvent.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/event/SnmpTransportListChangedEventListener.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/ProtoTransportEntityService.java
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportBalancingService.java
rename common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/{session/SnmpSessionListener.java => service/SnmpTransportService.java} (60%)
create mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java
delete mode 100644 common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionCtx.java
create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/DeviceUpdatedEvent.java
create mode 100644 docker/tb-transports/snmp/conf/logback.xml
create mode 100644 docker/tb-transports/snmp/conf/tb-snmp-transport.conf
create mode 100644 transport/snmp/src/main/conf/logback.xml
create mode 100644 transport/snmp/src/main/conf/tb-snmp-transport.conf
create mode 100644 transport/snmp/src/main/resources/logback.xml
create mode 100644 transport/snmp/src/main/resources/tb-snmp-transport.yml
diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
index 05363dfd59..f3a07eb727 100644
--- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
+++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
@@ -25,7 +25,6 @@ import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.DefaultTbActorSystem;
-import org.thingsboard.server.actors.TbActorId;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.actors.TbActorSystem;
import org.thingsboard.server.actors.TbActorSystemSettings;
@@ -33,14 +32,13 @@ import org.thingsboard.server.actors.app.AppActor;
import org.thingsboard.server.actors.app.AppInitMsg;
import org.thingsboard.server.actors.stats.StatsActor;
import org.thingsboard.server.common.msg.queue.PartitionChangeMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
@Service
@Slf4j
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
index d0a3984660..ccac9b4a16 100644
--- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
@@ -52,7 +52,7 @@ import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
index 117607ae9f..f97188597b 100644
--- a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
@@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
public interface TbApiUsageStateService extends ApplicationListener {
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
index 02f69ce10d..9d3e2b5a64 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
@@ -52,7 +52,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
index 390798a3e2..bd29b6c431 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
@@ -38,7 +38,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerService.java
index 0a965413cf..222f58fd78 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerService.java
@@ -16,7 +16,7 @@
package org.thingsboard.server.service.queue;
import org.springframework.context.ApplicationListener;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
public interface TbCoreConsumerService extends ApplicationListener {
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/TbRuleEngineConsumerService.java
index e4a88de575..582dbd1355 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/TbRuleEngineConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/TbRuleEngineConsumerService.java
@@ -16,7 +16,7 @@
package org.thingsboard.server.service.queue;
import org.springframework.context.ApplicationListener;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
public interface TbRuleEngineConsumerService extends ApplicationListener {
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
index 02378eb557..28cce50bf5 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
@@ -34,7 +34,7 @@ import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
index c8a4e80e0e..5c358c1b77 100644
--- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
@@ -54,7 +54,7 @@ import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.util.TbCoreComponent;
diff --git a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
index 932ed4720d..b1bbb9cdfa 100644
--- a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
@@ -18,7 +18,7 @@ package org.thingsboard.server.service.state;
import org.springframework.context.ApplicationListener;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.common.msg.queue.TbCallback;
diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java
index ce8838aaa8..657bdf6f75 100644
--- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java
@@ -46,7 +46,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionUpdate
import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionUpdateValueListProto;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
index 0220a94964..710d613b37 100644
--- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
@@ -20,10 +20,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
-import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.queue.discovery.ClusterTopologyChangeEvent;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java
index fd40364fe5..37850b701f 100644
--- a/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionManagerService.java
@@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.queue.TbCallback;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import java.util.List;
diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbLocalSubscriptionService.java
index d244ef136b..f1bf6c0da0 100644
--- a/application/src/main/java/org/thingsboard/server/service/subscription/TbLocalSubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbLocalSubscriptionService.java
@@ -15,8 +15,8 @@
*/
package org.thingsboard.server.service.subscription;
-import org.thingsboard.server.queue.discovery.ClusterTopologyChangeEvent;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate;
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate;
diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java
index 168e1271fd..23abba15f3 100644
--- a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java
@@ -22,35 +22,18 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationListener;
import org.springframework.context.event.EventListener;
-import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
-import org.thingsboard.server.common.data.id.EntityId;
-import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.kv.AttributeKvEntry;
-import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
-import org.thingsboard.server.common.data.kv.BooleanDataEntry;
-import org.thingsboard.server.common.data.kv.DoubleDataEntry;
-import org.thingsboard.server.common.data.kv.LongDataEntry;
-import org.thingsboard.server.common.data.kv.StringDataEntry;
-import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
-import org.thingsboard.server.dao.attributes.AttributesService;
-import org.thingsboard.server.dao.timeseries.TimeseriesService;
-import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.subscription.SubscriptionManagerService;
-import org.thingsboard.server.service.subscription.TbSubscriptionUtils;
import javax.annotation.Nullable;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
-import java.util.Collections;
-import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/AlarmSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/AlarmSubscriptionService.java
index fa44906cc3..6ed3860341 100644
--- a/application/src/main/java/org/thingsboard/server/service/telemetry/AlarmSubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/telemetry/AlarmSubscriptionService.java
@@ -17,8 +17,7 @@ package org.thingsboard.server.service.telemetry;
import org.springframework.context.ApplicationListener;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
-import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
/**
* Created by ashvayka on 27.03.18.
diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java
index 107b97dbb3..dde71ef53f 100644
--- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java
@@ -22,9 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
-import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmQuery;
@@ -35,43 +33,22 @@ import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.kv.AttributeKvEntry;
-import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
-import org.thingsboard.server.common.data.kv.BooleanDataEntry;
-import org.thingsboard.server.common.data.kv.DoubleDataEntry;
-import org.thingsboard.server.common.data.kv.LongDataEntry;
-import org.thingsboard.server.common.data.kv.StringDataEntry;
-import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.AlarmData;
-import org.thingsboard.server.common.data.query.AlarmDataPageLink;
import org.thingsboard.server.common.data.query.AlarmDataQuery;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.alarm.AlarmOperationResult;
import org.thingsboard.server.dao.alarm.AlarmService;
-import org.thingsboard.server.dao.attributes.AttributesService;
-import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.subscription.SubscriptionManagerService;
import org.thingsboard.server.service.subscription.TbSubscriptionUtils;
-import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate;
-import javax.annotation.PostConstruct;
-import javax.annotation.PreDestroy;
import java.util.Collection;
-import java.util.Collections;
-import java.util.List;
import java.util.Optional;
-import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-import java.util.function.Consumer;
/**
* Created by ashvayka on 27.03.18.
diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetrySubscriptionService.java
index 015a58495f..0ac9e79ff1 100644
--- a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetrySubscriptionService.java
+++ b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetrySubscriptionService.java
@@ -16,8 +16,7 @@
package org.thingsboard.server.service.telemetry;
import org.springframework.context.ApplicationListener;
-import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
-import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
/**
* Created by ashvayka on 27.03.18.
diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
index eea7af582c..4f58bd0f63 100644
--- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
+++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
@@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.ApiUsageState;
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;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
@@ -61,11 +62,15 @@ import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
+import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceCredentialsRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetEntityProfileResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetResourcesRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.GetSnmpDevicesRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.GetSnmpDevicesResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionResponseStatus;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
@@ -84,6 +89,7 @@ import org.thingsboard.server.service.state.DeviceStateService;
import java.util.Collections;
import java.util.List;
+import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -139,40 +145,41 @@ public class DefaultTransportApiService implements TransportApiService {
@Override
public ListenableFuture> handle(TbProtoQueueMsg tbProtoQueueMsg) {
TransportApiRequestMsg transportApiRequestMsg = tbProtoQueueMsg.getValue();
+ ListenableFuture result = null;
+
if (transportApiRequestMsg.hasValidateTokenRequestMsg()) {
ValidateDeviceTokenRequestMsg msg = transportApiRequestMsg.getValidateTokenRequestMsg();
- return Futures.transform(validateCredentials(msg.getToken(), DeviceCredentialsType.ACCESS_TOKEN),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = validateCredentials(msg.getToken(), DeviceCredentialsType.ACCESS_TOKEN);
} else if (transportApiRequestMsg.hasValidateBasicMqttCredRequestMsg()) {
TransportProtos.ValidateBasicMqttCredRequestMsg msg = transportApiRequestMsg.getValidateBasicMqttCredRequestMsg();
- return Futures.transform(validateCredentials(msg),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = validateCredentials(msg);
} else if (transportApiRequestMsg.hasValidateX509CertRequestMsg()) {
ValidateDeviceX509CertRequestMsg msg = transportApiRequestMsg.getValidateX509CertRequestMsg();
- return Futures.transform(validateCredentials(msg.getHash(), DeviceCredentialsType.X509_CERTIFICATE),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = validateCredentials(msg.getHash(), DeviceCredentialsType.X509_CERTIFICATE);
} else if (transportApiRequestMsg.hasGetOrCreateDeviceRequestMsg()) {
- return Futures.transform(handle(transportApiRequestMsg.getGetOrCreateDeviceRequestMsg()),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = handle(transportApiRequestMsg.getGetOrCreateDeviceRequestMsg());
} else if (transportApiRequestMsg.hasEntityProfileRequestMsg()) {
- return Futures.transform(handle(transportApiRequestMsg.getEntityProfileRequestMsg()),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = handle(transportApiRequestMsg.getEntityProfileRequestMsg());
} else if (transportApiRequestMsg.hasLwM2MRequestMsg()) {
- return Futures.transform(handle(transportApiRequestMsg.getLwM2MRequestMsg()),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = handle(transportApiRequestMsg.getLwM2MRequestMsg());
} else if (transportApiRequestMsg.hasValidateDeviceLwM2MCredentialsRequestMsg()) {
ValidateDeviceLwM2MCredentialsRequestMsg msg = transportApiRequestMsg.getValidateDeviceLwM2MCredentialsRequestMsg();
- return Futures.transform(validateCredentials(msg.getCredentialsId(), DeviceCredentialsType.LWM2M_CREDENTIALS),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = validateCredentials(msg.getCredentialsId(), DeviceCredentialsType.LWM2M_CREDENTIALS);
} else if (transportApiRequestMsg.hasProvisionDeviceRequestMsg()) {
- return Futures.transform(handle(transportApiRequestMsg.getProvisionDeviceRequestMsg()),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = handle(transportApiRequestMsg.getProvisionDeviceRequestMsg());
} else if (transportApiRequestMsg.hasResourcesRequestMsg()) {
- return Futures.transform(handle(transportApiRequestMsg.getResourcesRequestMsg()),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+ result = handle(transportApiRequestMsg.getResourcesRequestMsg());
+ } else if (transportApiRequestMsg.hasSnmpDevicesRequestMsg()) {
+ result = handle(transportApiRequestMsg.getSnmpDevicesRequestMsg());
+ } else if (transportApiRequestMsg.hasDeviceRequestMsg()) {
+ result = handle(transportApiRequestMsg.getDeviceRequestMsg());
+ } else if (transportApiRequestMsg.hasDeviceCredentialsRequestMsg()) {
+ result = handle(transportApiRequestMsg.getDeviceCredentialsRequestMsg());
}
- return Futures.transform(getEmptyTransportApiResponseFuture(),
- value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
+
+ return Futures.transform(Optional.ofNullable(result).orElseGet(this::getEmptyTransportApiResponseFuture),
+ value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()),
+ MoreExecutors.directExecutor());
}
private ListenableFuture validateCredentials(String credentialsId, DeviceCredentialsType credentialsType) {
@@ -366,6 +373,39 @@ public class DefaultTransportApiService implements TransportApiService {
return Futures.immediateFuture(TransportApiResponseMsg.newBuilder().setEntityProfileResponseMsg(builder).build());
}
+ private ListenableFuture handle(GetDeviceRequestMsg requestMsg) {
+ DeviceId deviceId = new DeviceId(new UUID(requestMsg.getDeviceIdMSB(), requestMsg.getDeviceIdLSB()));
+ Device device = deviceService.findDeviceById(TenantId.SYS_TENANT_ID, deviceId);
+
+ TransportApiResponseMsg responseMsg;
+ if (device != null) {
+ UUID deviceProfileId = device.getDeviceProfileId().getId();
+ responseMsg = TransportApiResponseMsg.newBuilder()
+ .setDeviceResponseMsg(TransportProtos.GetDeviceResponseMsg.newBuilder()
+ .setDeviceProfileIdMSB(deviceProfileId.getMostSignificantBits())
+ .setDeviceProfileIdLSB(deviceProfileId.getLeastSignificantBits())
+ .setDeviceTransportConfiguration(ByteString.copyFrom(
+ dataDecodingEncodingService.encode(device.getDeviceData().getTransportConfiguration())
+ )))
+ .build();
+ } else {
+ responseMsg = TransportApiResponseMsg.getDefaultInstance();
+ }
+
+ return Futures.immediateFuture(responseMsg);
+ }
+
+ private ListenableFuture handle(GetDeviceCredentialsRequestMsg requestMsg) {
+ DeviceId deviceId = new DeviceId(new UUID(requestMsg.getDeviceIdMSB(), requestMsg.getDeviceIdLSB()));
+ DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(TenantId.SYS_TENANT_ID, deviceId);
+
+ return Futures.immediateFuture(TransportApiResponseMsg.newBuilder()
+ .setDeviceCredentialsResponseMsg(TransportProtos.GetDeviceCredentialsResponseMsg.newBuilder()
+ .setDeviceCredentialsData(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceCredentials))))
+ .build());
+ }
+
+
private ListenableFuture handle(GetResourcesRequestMsg requestMsg) {
TenantId tenantId = new TenantId(new UUID(requestMsg.getTenantIdMSB(), requestMsg.getTenantIdLSB()));
TransportProtos.GetResourcesResponseMsg.Builder builder = TransportProtos.GetResourcesResponseMsg.newBuilder();
@@ -400,6 +440,20 @@ public class DefaultTransportApiService implements TransportApiService {
.build();
}
+ // TODO: request snmp devices with pagination
+ private ListenableFuture handle(GetSnmpDevicesRequestMsg requestMsg) {
+ List result = deviceService.findDevicesIdsByDeviceProfileTransportType(DeviceTransportType.SNMP);
+ GetSnmpDevicesResponseMsg responseMsg = GetSnmpDevicesResponseMsg.newBuilder()
+ .addAllIds(result.stream()
+ .map(UUID::toString)
+ .collect(Collectors.toList()))
+ .build();
+
+ return Futures.immediateFuture(TransportApiResponseMsg.newBuilder()
+ .setSnmpDevicesResponseMsg(responseMsg)
+ .build());
+ }
+
private ListenableFuture getDeviceInfo(DeviceId deviceId, DeviceCredentials credentials) {
return Futures.transform(deviceService.findDeviceByIdAsync(TenantId.SYS_TENANT_ID, deviceId), device -> {
if (device == null) {
diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
index f98c56d89a..1a893f5388 100644
--- a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
+++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
@@ -19,6 +19,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceInfo;
import org.thingsboard.server.common.data.DeviceProfile;
+import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.common.data.id.CustomerId;
@@ -31,6 +32,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.dao.device.provision.ProvisionRequest;
import java.util.List;
+import java.util.UUID;
public interface DeviceService {
@@ -90,4 +92,5 @@ public interface DeviceService {
Device saveDevice(ProvisionRequest provisionRequest, DeviceProfile profile);
+ List findDevicesIdsByDeviceProfileTransportType(DeviceTransportType transportType);
}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/TbTransportService.java b/common/data/src/main/java/org/thingsboard/server/common/data/TbTransportService.java
new file mode 100644
index 0000000000..6195022aba
--- /dev/null
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/TbTransportService.java
@@ -0,0 +1,20 @@
+/**
+ * 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.common.data;
+
+public interface TbTransportService {
+ String getName();
+}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceTransportConfiguration.java
index 0ae9f26cf3..49a547371f 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceTransportConfiguration.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceTransportConfiguration.java
@@ -21,6 +21,8 @@ import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.server.common.data.DeviceTransportType;
+import java.io.Serializable;
+
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
@@ -31,7 +33,7 @@ import org.thingsboard.server.common.data.DeviceTransportType;
@JsonSubTypes.Type(value = MqttDeviceTransportConfiguration.class, name = "MQTT"),
@JsonSubTypes.Type(value = Lwm2mDeviceTransportConfiguration.class, name = "LWM2M"),
@JsonSubTypes.Type(value = SnmpDeviceTransportConfiguration.class, name = "SNMP")})
-public interface DeviceTransportConfiguration {
+public interface DeviceTransportConfiguration extends Serializable {
@JsonIgnore
DeviceTransportType getType();
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpProfileTransportConfiguration.java
index a2529e74e9..50f6c4351c 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpProfileTransportConfiguration.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpProfileTransportConfiguration.java
@@ -19,14 +19,14 @@ import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.Data;
import org.thingsboard.server.common.data.DeviceTransportType;
+import java.util.Collections;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@Data
public class SnmpProfileTransportConfiguration implements DeviceProfileTransportConfiguration {
-
- private int poolPeriodMs;
+ private int pollPeriodMs;
private int timeoutMs;
private int retries;
private List attributes;
@@ -39,6 +39,10 @@ public class SnmpProfileTransportConfiguration implements DeviceProfileTransport
@JsonIgnore
public List getKvMappings() {
- return Stream.concat(attributes.stream(), telemetry.stream()).collect(Collectors.toList());
+ if (attributes != null && telemetry != null) {
+ return Stream.concat(attributes.stream(), telemetry.stream()).collect(Collectors.toList());
+ } else {
+ return Collections.emptyList();
+ }
}
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java
index fc4ec33eb1..3ba360cd8e 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java
@@ -19,8 +19,12 @@ import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.event.ContextRefreshedEvent;
+import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
+import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.transport.TransportProtos;
@@ -32,6 +36,7 @@ import javax.annotation.PostConstruct;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Arrays;
+import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
@@ -56,6 +61,8 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
@Autowired(required = false)
private TbQueueRuleEngineSettings ruleEngineSettings;
+ @Autowired
+ private ApplicationContext applicationContext;
private List serviceTypes;
private ServiceInfo serviceInfo;
@@ -102,6 +109,19 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
serviceInfo = builder.build();
}
+ @EventListener(ContextRefreshedEvent.class)
+ public void setTransports() {
+ serviceInfo = ServiceInfo.newBuilder(serviceInfo)
+ .addAllTransports(getTransportServices().stream()
+ .map(TbTransportService::getName)
+ .collect(Collectors.toSet()))
+ .build();
+ }
+
+ private Collection getTransportServices() {
+ return applicationContext.getBeansOfType(TbTransportService.class).values();
+ }
+
@Override
public ServiceInfo getServiceInfo() {
return serviceInfo;
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 2da438417a..c833b314e6 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
@@ -15,26 +15,27 @@
*/
package org.thingsboard.server.queue.discovery;
-import com.google.common.hash.HashCode;
import com.google.common.hash.HashFunction;
+import com.google.common.hash.Hasher;
import com.google.common.hash.Hashing;
-import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.msg.queue.ServiceQueueKey;
import org.thingsboard.server.common.msg.queue.ServiceQueue;
+import org.thingsboard.server.common.msg.queue.ServiceQueueKey;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
+import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
+import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent;
import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
import javax.annotation.PostConstruct;
-import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
@@ -46,7 +47,6 @@ import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
-import java.util.concurrent.ConcurrentNavigableMap;
import java.util.stream.Collectors;
@Service
@@ -186,6 +186,8 @@ public class HashPartitionService implements PartitionService {
applicationEventPublisher.publishEvent(new ClusterTopologyChangeEvent(this, changes));
}
}
+
+ applicationEventPublisher.publishEvent(new ServiceListChangedEvent(otherServices, currentService));
}
@Override
@@ -219,6 +221,14 @@ public class HashPartitionService implements PartitionService {
}
}
+ @Override
+ public int resolvePartitionIndex(UUID entityId, int partitions) {
+ int hash = hashFunction.newHasher()
+ .putLong(entityId.getMostSignificantBits())
+ .putLong(entityId.getLeastSignificantBits()).hash().asInt();
+ return Math.abs(hash % partitions);
+ }
+
private Map> getServiceKeyListMap(List services) {
final Map> currentMap = new HashMap<>();
services.forEach(serviceInfo -> {
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 ebb260bc57..20c59378e8 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
@@ -20,9 +20,11 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
+import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import java.util.List;
import java.util.Set;
+import java.util.UUID;
/**
* Once application is ready or cluster topology changes, this Service will produce {@link PartitionChangeEvent}
@@ -55,4 +57,6 @@ public interface PartitionService {
* @return
*/
TopicPartitionInfo getNotificationsTopic(ServiceType serviceType, String serviceId);
+
+ int resolvePartitionIndex(UUID entityId, int partitions);
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEventListener.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEventListener.java
index 9158d8f0c8..70484a892b 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEventListener.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEventListener.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.queue.discovery;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.ApplicationListener;
+import org.thingsboard.server.queue.discovery.event.TbApplicationEvent;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
index 5253f57248..0802f0f52c 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
@@ -33,12 +33,14 @@ import org.apache.zookeeper.KeeperException;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.event.EventListener;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Service;
import org.springframework.util.Assert;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.gen.transport.TransportProtos;
+import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -77,7 +79,8 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
private volatile boolean stopped = true;
- public ZkDiscoveryService(TbServiceInfoProvider serviceInfoProvider, PartitionService partitionService) {
+ public ZkDiscoveryService(TbServiceInfoProvider serviceInfoProvider,
+ PartitionService partitionService) {
this.serviceInfoProvider = serviceInfoProvider;
this.partitionService = partitionService;
}
@@ -126,7 +129,8 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
return;
}
publishCurrentServer();
- partitionService.recalculatePartitions(serviceInfoProvider.getServiceInfo(), getOtherServers());
+ TransportProtos.ServiceInfo currentService = serviceInfoProvider.getServiceInfo();
+ partitionService.recalculatePartitions(currentService, getOtherServers());
}
public synchronized void publishCurrentServer() {
@@ -281,11 +285,11 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
case CHILD_ADDED:
case CHILD_UPDATED:
case CHILD_REMOVED:
- partitionService.recalculatePartitions(serviceInfoProvider.getServiceInfo(), getOtherServers());
+ TransportProtos.ServiceInfo currentService = serviceInfoProvider.getServiceInfo();
+ partitionService.recalculatePartitions(currentService, getOtherServers());
break;
default:
break;
}
}
-
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ClusterTopologyChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ClusterTopologyChangeEvent.java
similarity index 91%
rename from common/queue/src/main/java/org/thingsboard/server/queue/discovery/ClusterTopologyChangeEvent.java
rename to common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ClusterTopologyChangeEvent.java
index 1e5b90b5fe..6b3a3ff93f 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ClusterTopologyChangeEvent.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ClusterTopologyChangeEvent.java
@@ -13,10 +13,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.queue.discovery;
+package org.thingsboard.server.queue.discovery.event;
import lombok.Getter;
-import org.springframework.context.ApplicationEvent;
import org.thingsboard.server.common.msg.queue.ServiceQueueKey;
import java.util.Set;
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java
similarity index 93%
rename from common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionChangeEvent.java
rename to common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java
index 2edcd2ceca..e11e19db61 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionChangeEvent.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java
@@ -13,10 +13,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.queue.discovery;
+package org.thingsboard.server.queue.discovery.event;
import lombok.Getter;
-import org.springframework.context.ApplicationEvent;
import org.thingsboard.server.common.msg.queue.ServiceQueueKey;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ServiceListChangedEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ServiceListChangedEvent.java
new file mode 100644
index 0000000000..dd7d384640
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/ServiceListChangedEvent.java
@@ -0,0 +1,35 @@
+/**
+ * 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.queue.discovery.event;
+
+import lombok.Getter;
+import lombok.ToString;
+import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
+
+import java.util.List;
+
+@Getter
+@ToString
+public class ServiceListChangedEvent extends TbApplicationEvent {
+ private final List otherServices;
+ private final ServiceInfo currentService;
+
+ public ServiceListChangedEvent(List otherServices, ServiceInfo currentService) {
+ super(otherServices);
+ this.otherServices = otherServices;
+ this.currentService = currentService;
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/TbApplicationEvent.java
similarity index 95%
rename from common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEvent.java
rename to common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/TbApplicationEvent.java
index face2d36d6..1a75b4e4e4 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbApplicationEvent.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/TbApplicationEvent.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.queue.discovery;
+package org.thingsboard.server.queue.discovery.event;
import lombok.Getter;
import org.springframework.context.ApplicationEvent;
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/TbSnmpTransportComponent.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbSnmpTransportComponent.java
new file mode 100644
index 0000000000..36080afaab
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/TbSnmpTransportComponent.java
@@ -0,0 +1,29 @@
+/**
+ * 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.queue.util;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.snmp.enabled}'=='true')")
+@Retention(RetentionPolicy.RUNTIME)
+@Target({ElementType.TYPE, ElementType.METHOD})
+public @interface TbSnmpTransportComponent {
+}
diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto
index 42f6267d42..59fcdad36a 100644
--- a/common/queue/src/main/proto/queue.proto
+++ b/common/queue/src/main/proto/queue.proto
@@ -34,6 +34,7 @@ message ServiceInfo {
int64 tenantIdMSB = 3;
int64 tenantIdLSB = 4;
repeated QueueInfo ruleEngineQueues = 5;
+ repeated string transports = 6;
}
/**
@@ -250,6 +251,34 @@ message GetEntityProfileResponseMsg {
bytes apiState = 3;
}
+message GetDeviceRequestMsg {
+ int64 deviceIdMSB = 1;
+ int64 deviceIdLSB = 2;
+}
+
+message GetDeviceResponseMsg {
+ int64 deviceProfileIdMSB = 1;
+ int64 deviceProfileIdLSB = 2;
+ bytes deviceTransportConfiguration = 3;
+}
+
+message GetDeviceCredentialsRequestMsg {
+ int64 deviceIdMSB = 1;
+ int64 deviceIdLSB = 2;
+}
+
+message GetDeviceCredentialsResponseMsg {
+ bytes deviceCredentialsData = 1;
+}
+
+message GetSnmpDevicesRequestMsg {
+
+}
+
+message GetSnmpDevicesResponseMsg {
+ repeated string ids = 1;
+}
+
message EntityUpdateMsg {
string entityType = 1;
bytes data = 2;
@@ -559,6 +588,9 @@ message TransportApiRequestMsg {
ProvisionDeviceRequestMsg provisionDeviceRequestMsg = 7;
ValidateDeviceLwM2MCredentialsRequestMsg validateDeviceLwM2MCredentialsRequestMsg = 8;
GetResourcesRequestMsg resourcesRequestMsg = 9;
+ GetSnmpDevicesRequestMsg snmpDevicesRequestMsg = 10;
+ GetDeviceRequestMsg deviceRequestMsg = 11;
+ GetDeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 12;
}
/* Response from ThingsBoard Core Service to Transport Service */
@@ -567,8 +599,11 @@ message TransportApiResponseMsg {
GetOrCreateDeviceFromGatewayResponseMsg getOrCreateDeviceResponseMsg = 2;
GetEntityProfileResponseMsg entityProfileResponseMsg = 3;
ProvisionDeviceResponseMsg provisionDeviceResponseMsg = 4;
+ GetSnmpDevicesResponseMsg snmpDevicesResponseMsg = 5;
LwM2MResponseMsg lwM2MResponseMsg = 6;
GetResourcesResponseMsg resourcesResponseMsg = 7;
+ GetDeviceResponseMsg deviceResponseMsg = 8;
+ GetDeviceCredentialsResponseMsg deviceCredentialsResponseMsg = 9;
}
/* Messages that are handled by ThingsBoard Core Service */
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 8abc859b21..dd4a76dcc8 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
@@ -18,12 +18,12 @@ package org.thingsboard.server.transport.coap;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.core.CoapResource;
import org.eclipse.californium.core.CoapServer;
-
import org.eclipse.californium.core.network.CoapEndpoint;
import org.eclipse.californium.core.network.CoapEndpoint.Builder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service;
+import org.thingsboard.server.common.data.TbTransportService;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -34,7 +34,7 @@ import java.net.UnknownHostException;
@Service("CoapTransportService")
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.coap.enabled}'=='true')")
@Slf4j
-public class CoapTransportService {
+public class CoapTransportService implements TbTransportService {
private static final String V1 = "v1";
private static final String API = "api";
@@ -73,4 +73,9 @@ public class CoapTransportService {
this.server.destroy();
log.info("CoAP transport stopped!");
}
+
+ @Override
+ public String getName() {
+ return "COAP";
+ }
}
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 e222b690ef..011093149e 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
@@ -31,6 +31,7 @@ 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.DeviceTransportType;
+import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportContext;
@@ -41,7 +42,6 @@ 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.gen.transport.TransportProtos.AttributeUpdateNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
@@ -53,7 +53,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMs
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg;
import javax.servlet.http.HttpServletRequest;
@@ -69,7 +68,7 @@ import java.util.function.Consumer;
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.http.enabled}'=='true')")
@RequestMapping("/api/v1")
@Slf4j
-public class DeviceApiController {
+public class DeviceApiController implements TbTransportService {
@Autowired
private HttpTransportContext transportContext;
@@ -338,4 +337,9 @@ public class DeviceApiController {
.build(), TransportServiceCallback.EMPTY);
}
+ @Override
+ public String getName() {
+ return "HTTP";
+ }
+
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java
index 8d1aff37b5..db419eca9f 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java
@@ -20,12 +20,13 @@ import org.eclipse.leshan.core.response.ReadResponse;
import org.eclipse.leshan.server.registration.Registration;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
+import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Collection;
import java.util.Optional;
-public interface LwM2mTransportService {
+public interface LwM2mTransportService extends TbTransportService {
void onRegistered(Registration registration, Collection previousObsersations);
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java
index 548d750a3f..bae6a0af1d 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java
@@ -40,6 +40,7 @@ import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
+import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
@@ -1118,4 +1119,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
namesIsWritable.addAll(new HashSet<>(keyNamesIsWritable.values()));
return new ArrayList<>(namesIsWritable);
}
+
+ @Override
+ public String getName() {
+ return "LWM2M";
+ }
}
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 0a0c0cbe70..3e53dec227 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.TbTransportService;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -38,7 +39,7 @@ import javax.annotation.PreDestroy;
@Service("MqttTransportService")
@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.mqtt.enabled}'=='true')")
@Slf4j
-public class MqttTransportService {
+public class MqttTransportService implements TbTransportService {
@Value("${transport.mqtt.bind_address}")
private String host;
@@ -90,4 +91,9 @@ public class MqttTransportService {
}
log.info("MQTT transport stopped!");
}
+
+ @Override
+ public String getName() {
+ return "MQTT";
+ }
}
diff --git a/common/transport/snmp/pom.xml b/common/transport/snmp/pom.xml
index c7e22ecc80..ab83ab69a7 100644
--- a/common/transport/snmp/pom.xml
+++ b/common/transport/snmp/pom.xml
@@ -58,9 +58,5 @@
org.snmp4j
snmp4j
-
- org.thingsboard.common
- dao-api
-
diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulator.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulator.java
new file mode 100644
index 0000000000..7109ef4e3f
--- /dev/null
+++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulator.java
@@ -0,0 +1,81 @@
+/**
+ * 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.snmp;
+
+import org.snmp4j.CommunityTarget;
+import org.snmp4j.PDU;
+import org.snmp4j.Snmp;
+import org.snmp4j.Target;
+import org.snmp4j.event.ResponseEvent;
+import org.snmp4j.mp.SnmpConstants;
+import org.snmp4j.smi.GenericAddress;
+import org.snmp4j.smi.OID;
+import org.snmp4j.smi.OctetString;
+import org.snmp4j.smi.VariableBinding;
+import org.snmp4j.transport.DefaultUdpTransportMapping;
+import org.snmp4j.transport.UdpTransportMapping;
+
+import java.io.IOException;
+
+/**
+ * For testing purposes. Will be removed when the time comes
+ */
+public class SnmpDeviceSimulator {
+ private final Target target;
+ private final OID oid = new OID(".1.3.6.1.2.1.1.1.0");
+ private Snmp snmp;
+
+ public SnmpDeviceSimulator(int port) {
+ String address = "udp:127.0.0.1/" + port;
+
+ CommunityTarget target = new CommunityTarget();
+ target.setCommunity(new OctetString("public"));
+ target.setAddress(GenericAddress.parse(address));
+ target.setRetries(2);
+ target.setTimeout(1500);
+ target.setVersion(SnmpConstants.version2c);
+
+ this.target = target;
+ }
+
+ public static void main(String[] args) throws IOException {
+ SnmpDeviceSimulator deviceSimulator = new SnmpDeviceSimulator(161);
+
+ deviceSimulator.start();
+ String response = deviceSimulator.sendRequest(PDU.GET);
+
+ System.out.println(response);
+ }
+
+ public void start() throws IOException {
+ UdpTransportMapping transport = new DefaultUdpTransportMapping();
+ transport.addTransportListener((sourceTransport, incomingAddress, wholeMessage, tmStateReference) -> {
+ System.out.println();
+ });
+ snmp = new Snmp(transport);
+
+ transport.listen();
+ }
+
+ public String sendRequest(int pduType) throws IOException {
+ PDU pdu = new PDU();
+ pdu.add(new VariableBinding(oid));
+ pdu.setType(pduType);
+
+ ResponseEvent responseEvent = snmp.send(pdu, target);
+ return responseEvent.getResponse().get(0).getVariable().toString();
+ }
+}
diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
index 66f852a209..a9fa8380fa 100644
--- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
+++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
@@ -15,17 +15,19 @@
*/
package org.thingsboard.server.transport.snmp;
-import lombok.Getter;
+import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.snmp4j.PDU;
-import org.snmp4j.Snmp;
import org.snmp4j.smi.OID;
import org.snmp4j.smi.VariableBinding;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.stereotype.Service;
+import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.event.EventListener;
+import org.springframework.core.annotation.Order;
+import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
+import org.thingsboard.server.common.data.DeviceTransportType;
+import org.thingsboard.server.common.data.device.data.DeviceTransportConfiguration;
import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration;
import org.thingsboard.server.common.data.device.profile.SnmpDeviceProfileKvMapping;
import org.thingsboard.server.common.data.device.profile.SnmpProfileTransportConfiguration;
@@ -33,84 +35,180 @@ import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
+import org.thingsboard.server.common.transport.DeviceUpdatedEvent;
import org.thingsboard.server.common.transport.TransportContext;
-import org.thingsboard.server.dao.device.DeviceCredentialsService;
-import org.thingsboard.server.transport.snmp.session.DeviceSessionCtx;
+import org.thingsboard.server.common.transport.TransportDeviceProfileCache;
+import org.thingsboard.server.common.transport.TransportService;
+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.common.transport.session.DeviceAwareSessionContext;
+import org.thingsboard.server.gen.transport.TransportProtos;
+import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
+import org.thingsboard.server.queue.util.TbSnmpTransportComponent;
+import org.thingsboard.server.transport.snmp.service.ProtoTransportEntityService;
+import org.thingsboard.server.transport.snmp.service.SnmpTransportBalancingService;
+import org.thingsboard.server.transport.snmp.service.SnmpTransportService;
+import org.thingsboard.server.transport.snmp.session.DeviceSessionContext;
import java.util.ArrayList;
+import java.util.Collection;
import java.util.HashMap;
+import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.concurrent.ExecutorService;
import java.util.stream.Collectors;
-@Service("SnmpTransportContext")
-@ConditionalOnExpression("'${service.type:null}'=='tb-transport' || ('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true' && '${transport.snmp.enabled}'=='true')")
+@TbSnmpTransportComponent
+@Component
@Slf4j
+@RequiredArgsConstructor
public class SnmpTransportContext extends TransportContext {
- @Autowired
- DeviceCredentialsService deviceCredentialsService;
-
- @Autowired
- SnmpTransportService snmpTransportService;
-
- @Getter
- private final Map profileTransportConfig = new ConcurrentHashMap<>();
- @Getter
- private final Map> pdusPerProfile = new ConcurrentHashMap<>();
- @Getter
- private final Map deviceSessions = new ConcurrentHashMap<>();
-
- public Optional findAttributesMapping(DeviceProfileId deviceProfileId, OID responseOid) {
- if (profileTransportConfig.containsKey(deviceProfileId)) {
- return findMapping(responseOid, profileTransportConfig.get(deviceProfileId).getAttributes());
- }
- return Optional.empty();
+ private final SnmpTransportService snmpTransportService;
+ private final TransportDeviceProfileCache deviceProfileCache;
+ private final TransportService transportService;
+ private final ProtoTransportEntityService protoEntityService;
+ private final SnmpTransportBalancingService balancingService;
+
+ private final Map