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 4cd5c664fd..0d0e38a1b2 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
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
import akka.actor.ActorRef;
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 96b72843ed..413d5ee9dd 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
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
import akka.actor.ActorRef;
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/MsgPackCallback.java b/application/src/main/java/org/thingsboard/server/service/queue/MsgPackCallback.java
index c700c18dc0..5b94db28b0 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/MsgPackCallback.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/MsgPackCallback.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
import lombok.extern.slf4j.Slf4j;
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 0541de7c63..1e1b352645 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
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
import org.thingsboard.server.gen.transport.TransportProtos;
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java
index e3bb3a3c4b..f7912836c5 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java
@@ -5,7 +5,7 @@
* 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
+ * 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,
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgCallback.java b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgCallback.java
index 44dd058e1e..774377407a 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgCallback.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgCallback.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
public interface TbMsgCallback {
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 691a54aaab..671cd72262 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
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.queue;
public interface TbRuleEngineConsumerService {
diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java
index c1b0608d44..49031f9f53 100644
--- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java
+++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.transport;
import lombok.extern.slf4j.Slf4j;
diff --git a/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java
index 354252642f..4765737430 100644
--- a/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java
+++ b/application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java
@@ -5,7 +5,7 @@
* 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
+ * 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,
diff --git a/application/src/main/java/org/thingsboard/server/service/transport/TbCoreToTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/TbCoreToTransportService.java
index 3324f02440..5b11edb5f0 100644
--- a/application/src/main/java/org/thingsboard/server/service/transport/TbCoreToTransportService.java
+++ b/application/src/main/java/org/thingsboard/server/service/transport/TbCoreToTransportService.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.service.transport;
import org.thingsboard.server.gen.transport.TransportProtos;
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 2ded608203..da2e5beb98 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -70,11 +70,6 @@ zk:
# Name of the directory in zookeeper 'filesystem'
zk_dir: "${ZOOKEEPER_NODES_DIR:/thingsboard}"
-# RPC connection parameters. Used only in cluster mode only.
-rpc:
- bind_host: "${RPC_HOST:localhost}"
- bind_port: "${RPC_PORT:9001}"
-
# Clustering properties related to consistent-hashing. See architecture docs for more details.
cluster:
# Unique id for this node (autogenerated if empty)
@@ -521,7 +516,7 @@ swagger:
api_path_regex: "${SWAGGER_API_PATH_REGEX:/api.*}"
security_path_regex: "${SWAGGER_SECURITY_PATH_REGEX:/api.*}"
non_security_path_regex: "${SWAGGER_NON_SECURITY_PATH_REGEX:/api/noauth.*}"
- title: "${SWAGGER_TITLE:Thingsboard REST API}"
+ title: "${SWAGGER_TITLE:ThingsBoard REST API}"
description: "${SWAGGER_DESCRIPTION:For instructions how to authorize requests please visit REST API documentation page.}"
contact:
name: "${SWAGGER_CONTACT_NAME:Thingsboard team}"
@@ -568,4 +563,7 @@ queue:
topic: "${TB_QUEUE_TRANSPORT_NOTIFICATIONS_TOPIC:tb.transport.notifications}"
service:
- type: "${TB_SERVICE_TYPE:monolith}" # monolith or tb-core or tb-rule-engine or tb-transport
\ No newline at end of file
+ type: "${TB_SERVICE_TYPE:monolith}" # monolith or tb-core or tb-rule-engine or tb-transport
+ # Unique id for this service (autogenerated if empty)
+ id: "${TB_SERVICE_ID:}"
+ tenant_id: "${TB_SERVICE_TENANT_ID:}" # empty or specific tenant id.
\ No newline at end of file
diff --git a/common/queue/pom.xml b/common/queue/pom.xml
index fd26358cc6..429f68acb0 100644
--- a/common/queue/pom.xml
+++ b/common/queue/pom.xml
@@ -40,6 +40,10 @@
org.thingsboard.common
data
+
+ org.thingsboard.common
+ util
+
org.thingsboard.common
message
@@ -84,6 +88,14 @@
ch.qos.logback
logback-classic
+
+ com.google.protobuf
+ protobuf-java
+
+
+ org.apache.curator
+ curator-recipes
+
junit
junit
@@ -94,10 +106,6 @@
mockito-all
test
-
- com.google.protobuf
- protobuf-java
-
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueAdmin.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueAdmin.java
index 7f8cff5f22..6917c23719 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueAdmin.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueAdmin.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import com.google.common.util.concurrent.ListenableFuture;
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueCallback.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueCallback.java
index 3d7d791ae4..e823d2a5fc 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueCallback.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueCallback.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
public interface TbQueueCallback {
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueConsumer.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueConsumer.java
index ca51eb4ddb..ddf9d7d9b3 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueConsumer.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueConsumer.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import java.util.List;
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsg.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsg.java
index 22714af77e..ca8f3a61da 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsg.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsg.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import java.util.UUID;
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgHeaders.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgHeaders.java
index 9dad87588d..f95454c976 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgHeaders.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgHeaders.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import java.util.Map;
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgMetadata.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgMetadata.java
index 8eecae51fb..a3331999a5 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgMetadata.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueMsgMetadata.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
public interface TbQueueMsgMetadata {
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueProducer.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueProducer.java
index 41d165d743..2ec01357b7 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueProducer.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueProducer.java
@@ -1,6 +1,22 @@
+/**
+ * Copyright © 2016-2020 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;
import com.google.common.util.concurrent.ListenableFuture;
+import org.thingsboard.server.discovery.TopicPartitionInfo;
public interface TbQueueProducer {
@@ -12,4 +28,6 @@ public interface TbQueueProducer {
ListenableFuture send(String topic, T msg, TbQueueCallback callback);
+ ListenableFuture send(String topic, int partition, T msg, TbQueueCallback callback);
+
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueRequestTemplate.java
index 5ba28f22c5..5182bb7260 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueRequestTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueRequestTemplate.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import com.google.common.util.concurrent.ListenableFuture;
diff --git a/common/queue/src/main/java/org/thingsboard/server/TbQueueResponseTemplate.java b/common/queue/src/main/java/org/thingsboard/server/TbQueueResponseTemplate.java
index f3b4e4ad5d..686570ca32 100644
--- a/common/queue/src/main/java/org/thingsboard/server/TbQueueResponseTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/TbQueueResponseTemplate.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
public interface TbQueueResponseTemplate {
diff --git a/common/queue/src/main/java/org/thingsboard/server/common/AbstractTbQueueTemplate.java b/common/queue/src/main/java/org/thingsboard/server/common/AbstractTbQueueTemplate.java
index 2f406a9631..118cf5e0b7 100644
--- a/common/queue/src/main/java/org/thingsboard/server/common/AbstractTbQueueTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/common/AbstractTbQueueTemplate.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import java.nio.ByteBuffer;
diff --git a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueMsgHeaders.java b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueMsgHeaders.java
index e547a6c070..de0c368a65 100644
--- a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueMsgHeaders.java
+++ b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueMsgHeaders.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import org.thingsboard.server.TbQueueMsgHeaders;
diff --git a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueRequestTemplate.java
index 49022ad2e3..6876c2f3d2 100644
--- a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueRequestTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueRequestTemplate.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import com.google.common.util.concurrent.Futures;
diff --git a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueResponseTemplate.java b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueResponseTemplate.java
index 92dac0551b..ba32b399ce 100644
--- a/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueResponseTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/common/DefaultTbQueueResponseTemplate.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import lombok.Builder;
diff --git a/common/queue/src/main/java/org/thingsboard/server/common/TbProtoQueueMsg.java b/common/queue/src/main/java/org/thingsboard/server/common/TbProtoQueueMsg.java
index c072dc4d67..8dd492504c 100644
--- a/common/queue/src/main/java/org/thingsboard/server/common/TbProtoQueueMsg.java
+++ b/common/queue/src/main/java/org/thingsboard/server/common/TbProtoQueueMsg.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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;
import lombok.Data;
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/DefaultTbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/discovery/DefaultTbServiceInfoProvider.java
new file mode 100644
index 0000000000..0b09a333d8
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/DefaultTbServiceInfoProvider.java
@@ -0,0 +1,90 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import lombok.Getter;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
+
+import javax.annotation.PostConstruct;
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.UUID;
+import java.util.stream.Collectors;
+
+@Component
+@Slf4j
+public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
+
+ @Getter
+ @Value("${service.id:#{null}}")
+ private String serviceId;
+
+ @Getter
+ @Value("${service.type:monolith}")
+ private String serviceType;
+
+ @Getter
+ @Value("${service.tenant_id:}")
+ private String tenantIdStr;
+
+ private List serviceTypes;
+ private ServiceInfo serviceInfo;
+
+ @PostConstruct
+ public void init() {
+ if (StringUtils.isEmpty(serviceId)) {
+ try {
+ serviceId = InetAddress.getLocalHost().getHostName();
+ } catch (UnknownHostException e) {
+ serviceId = org.apache.commons.lang3.RandomStringUtils.randomAlphabetic(10);
+ }
+ }
+ log.info("Current Service ID: {}", serviceId);
+ if (serviceType.equalsIgnoreCase("monolith")) {
+ serviceTypes = Collections.unmodifiableList(Arrays.asList(ServiceType.values()));
+ } else {
+ serviceTypes = Collections.singletonList(ServiceType.valueOf(serviceType));
+ }
+ ServiceInfo.Builder builder = ServiceInfo.newBuilder()
+ .setServiceId(serviceId)
+ .addAllServiceTypes(serviceTypes.stream().map(ServiceType::name).collect(Collectors.toList()));
+ if (!StringUtils.isEmpty(tenantIdStr)) {
+ UUID tenantId = UUID.fromString(tenantIdStr);
+ builder.setTenantIdMSB(tenantId.getMostSignificantBits());
+ builder.setTenantIdLSB(tenantId.getLeastSignificantBits());
+ }
+ serviceInfo = builder.build();
+ }
+
+
+ @Override
+ public List getSupportedServiceTypes() {
+ return serviceTypes;
+ }
+
+ @Override
+ public ServiceInfo getServiceInfo() {
+ return serviceInfo;
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionChangeEvent.java
new file mode 100644
index 0000000000..7d083ba873
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionChangeEvent.java
@@ -0,0 +1,33 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import lombok.Getter;
+import org.springframework.context.ApplicationEvent;
+
+import java.util.List;
+
+
+public class PartitionChangeEvent extends ApplicationEvent {
+
+ @Getter
+ private final List partitions;
+
+ public PartitionChangeEvent(Object source, List partitions) {
+ super(source);
+ this.partitions = partitions;
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionDiscoveryService.java
new file mode 100644
index 0000000000..d85b1d3f44
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/PartitionDiscoveryService.java
@@ -0,0 +1,29 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
+
+import java.util.List;
+
+public interface PartitionDiscoveryService {
+
+ List getCurrentPartitions(ServiceType serviceType);
+
+ TopicPartitionInfo resolve(ServiceType serviceType, TenantId tenantId, EntityId entityId);
+
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/ServiceType.java b/common/queue/src/main/java/org/thingsboard/server/discovery/ServiceType.java
new file mode 100644
index 0000000000..801069eb61
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/ServiceType.java
@@ -0,0 +1,20 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+public enum ServiceType {
+ TB_CORE, TB_RULE_ENGINE, TB_TRANSPORT, JS_EXECUTOR
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/TbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/discovery/TbServiceInfoProvider.java
new file mode 100644
index 0000000000..5cb545c5d6
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/TbServiceInfoProvider.java
@@ -0,0 +1,30 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
+
+import java.util.List;
+
+public interface TbServiceInfoProvider {
+
+ List getSupportedServiceTypes();
+
+ String getServiceId();
+
+ ServiceInfo getServiceInfo();
+
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/TopicPartitionInfo.java b/common/queue/src/main/java/org/thingsboard/server/discovery/TopicPartitionInfo.java
new file mode 100644
index 0000000000..b46b941360
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/TopicPartitionInfo.java
@@ -0,0 +1,26 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import lombok.Data;
+
+@Data
+public class TopicPartitionInfo {
+
+ private String topic;
+ private int partition;
+
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/discovery/ZkPartitionDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/discovery/ZkPartitionDiscoveryService.java
new file mode 100644
index 0000000000..a722f2f395
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/discovery/ZkPartitionDiscoveryService.java
@@ -0,0 +1,255 @@
+/**
+ * Copyright © 2016-2020 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.discovery;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.SerializationUtils;
+import org.apache.curator.framework.CuratorFramework;
+import org.apache.curator.framework.CuratorFrameworkFactory;
+import org.apache.curator.framework.imps.CuratorFrameworkState;
+import org.apache.curator.framework.recipes.cache.PathChildrenCache;
+import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent;
+import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener;
+import org.apache.curator.framework.state.ConnectionState;
+import org.apache.curator.framework.state.ConnectionStateListener;
+import org.apache.curator.retry.RetryForever;
+import org.apache.curator.utils.CloseableUtils;
+import org.apache.zookeeper.CreateMode;
+import org.apache.zookeeper.KeeperException;
+import org.springframework.beans.factory.annotation.Autowired;
+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.event.EventListener;
+import org.springframework.stereotype.Service;
+import org.springframework.util.Assert;
+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.msg.cluster.ServerAddress;
+
+import javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+@Service
+@ConditionalOnProperty(prefix = "zk", value = "enabled", havingValue = "true", matchIfMissing = false)
+@Slf4j
+public class ZkPartitionDiscoveryService implements PartitionDiscoveryService, PathChildrenCacheListener {
+
+ @Value("${zk.url}")
+ private String zkUrl;
+ @Value("${zk.retry_interval_ms}")
+ private Integer zkRetryInterval;
+ @Value("${zk.connection_timeout_ms}")
+ private Integer zkConnectionTimeout;
+ @Value("${zk.session_timeout_ms}")
+ private Integer zkSessionTimeout;
+ @Value("${zk.zk_dir}")
+ private String zkDir;
+
+ @Value("${queue.core.partitions:100}")
+ private Integer corePartitions;
+ @Value("${queue.rule_engine.partitions:100}")
+ private Integer ruleEnginePartitions;
+
+ @Autowired
+ private TbServiceInfoProvider serviceIdProvider;
+
+ private final ConcurrentMap partitionSizes = new ConcurrentHashMap<>();
+ private final ConcurrentMap> myPartitions = new ConcurrentHashMap<>();
+
+ private ExecutorService reconnectExecutorService;
+ private CuratorFramework client;
+ private PathChildrenCache cache;
+ private String nodePath;
+ private String zkNodesDir;
+
+ private volatile boolean stopped = true;
+
+ @Override
+ public List getCurrentPartitions(ServiceType serviceType) {
+ return Collections.emptyList();
+ }
+
+ @Override
+ public TopicPartitionInfo resolve(ServiceType serviceType, TenantId tenantId, EntityId entityId) {
+
+ }
+
+ @PostConstruct
+ public void init() {
+ log.info("Initializing...");
+ Assert.hasLength(zkUrl, missingProperty("zk.url"));
+ Assert.notNull(zkRetryInterval, missingProperty("zk.retry_interval_ms"));
+ Assert.notNull(zkConnectionTimeout, missingProperty("zk.connection_timeout_ms"));
+ Assert.notNull(zkSessionTimeout, missingProperty("zk.session_timeout_ms"));
+
+ reconnectExecutorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("zk-discovery"));
+
+ partitionSizes.put(ServiceType.TB_CORE, corePartitions);
+ partitionSizes.put(ServiceType.TB_RULE_ENGINE, ruleEnginePartitions);
+
+ log.info("Initializing discovery service using ZK connect string: {}", zkUrl);
+
+ zkNodesDir = zkDir + "/nodes";
+ initZkClient();
+ }
+
+ @EventListener(ApplicationReadyEvent.class)
+ public void onApplicationEvent(ApplicationReadyEvent event) {
+ if (stopped) {
+ log.debug("Ignoring application ready event. Service is stopped.");
+ return;
+ } else {
+ log.info("Received application ready event. Starting current ZK node.");
+ }
+ if (client.getState() != CuratorFrameworkState.STARTED) {
+ log.debug("Ignoring application ready event, ZK client is not started, ZK client state [{}]", client.getState());
+ return;
+ }
+ publishCurrentServer();
+ getOtherServers().forEach(
+ server -> log.info("Found active server: [{}:{}]", server.getHost(), server.getPort())
+ );
+ }
+
+ @Override
+ public synchronized void publishCurrentServer() {
+ ServerInstance self = this.serverInstance.getSelf();
+ if (currentServerExists()) {
+ log.info("[{}:{}] ZK node for current instance already exists, NOT created new one: {}", self.getHost(), self.getPort(), nodePath);
+ } else {
+ try {
+ log.info("[{}:{}] Creating ZK node for current instance", self.getHost(), self.getPort());
+ nodePath = client.create()
+ .creatingParentsIfNeeded()
+ .withMode(CreateMode.EPHEMERAL_SEQUENTIAL).forPath(zkNodesDir + "/", SerializationUtils.serialize(self.getServerAddress()));
+ log.info("[{}:{}] Created ZK node for current instance: {}", self.getHost(), self.getPort(), nodePath);
+ client.getConnectionStateListenable().addListener(checkReconnect(self));
+ } catch (Exception e) {
+ log.error("Failed to create ZK node", e);
+ throw new RuntimeException(e);
+ }
+ }
+ }
+
+ private boolean currentServerExists() {
+ if (nodePath == null) {
+ return false;
+ }
+ try {
+ ServerInstance self = this.serverInstance.getSelf();
+ ServerAddress registeredServerAdress = null;
+ registeredServerAdress = SerializationUtils.deserialize(client.getData().forPath(nodePath));
+ if (self.getServerAddress() != null && self.getServerAddress().equals(registeredServerAdress)) {
+ return true;
+ }
+ } catch (KeeperException.NoNodeException e) {
+ log.info("ZK node does not exist: {}", nodePath);
+ } catch (Exception e) {
+ log.error("Couldn't check if ZK node exists", e);
+ }
+ return false;
+ }
+
+ private ConnectionStateListener checkReconnect(ServerInstance self) {
+ return (client, newState) -> {
+ log.info("[{}:{}] ZK state changed: {}", self.getHost(), self.getPort(), newState);
+ if (newState == ConnectionState.LOST) {
+ reconnectExecutorService.submit(this::reconnect);
+ }
+ };
+ }
+
+ private volatile boolean reconnectInProgress = false;
+
+ private synchronized void reconnect() {
+ if (!reconnectInProgress) {
+ reconnectInProgress = true;
+ try {
+ destroyZkClient();
+ initZkClient();
+ publishCurrentServer();
+ } catch (Exception e) {
+ log.error("Failed to reconnect to ZK: {}", e.getMessage(), e);
+ } finally {
+ reconnectInProgress = false;
+ }
+ }
+ }
+
+ private void initZkClient() {
+ try {
+ client = CuratorFrameworkFactory.newClient(zkUrl, zkSessionTimeout, zkConnectionTimeout, new RetryForever(zkRetryInterval));
+ client.start();
+ client.blockUntilConnected();
+ cache = new PathChildrenCache(client, zkNodesDir, true);
+ cache.getListenable().addListener(this);
+ cache.start();
+ stopped = false;
+ log.info("ZK client connected");
+ } catch (Exception e) {
+ log.error("Failed to connect to ZK: {}", e.getMessage(), e);
+ CloseableUtils.closeQuietly(cache);
+ CloseableUtils.closeQuietly(client);
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void unpublishCurrentServer() {
+ try {
+ if (nodePath != null) {
+ client.delete().forPath(nodePath);
+ }
+ } catch (Exception e) {
+ log.error("Failed to delete ZK node {}", nodePath, e);
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void destroyZkClient() {
+ stopped = true;
+ try {
+ unpublishCurrentServer();
+ } catch (Exception e) {
+ }
+ CloseableUtils.closeQuietly(cache);
+ CloseableUtils.closeQuietly(client);
+ log.info("ZK client disconnected");
+ }
+
+ @PreDestroy
+ public void destroy() {
+ destroyZkClient();
+ reconnectExecutorService.shutdownNow();
+ log.info("Stopped discovery service");
+ }
+
+ public static String missingProperty(String propertyName) {
+ return "The " + propertyName + " property need to be set!";
+ }
+
+ @Override
+ public void childEvent(CuratorFramework curatorFramework, PathChildrenCacheEvent pathChildrenCacheEvent) throws Exception {
+
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsg.java b/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsg.java
index e44a93a70d..9f8e73bdfc 100644
--- a/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsg.java
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsg.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.kafka;
import org.apache.kafka.clients.consumer.ConsumerRecord;
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsgMetadata.java b/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsgMetadata.java
index 4698c226d6..09cc292deb 100644
--- a/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsgMetadata.java
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/KafkaTbQueueMsgMetadata.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.kafka;
import lombok.AllArgsConstructor;
diff --git a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryStorage.java
index abf30667b5..ded4cd9810 100644
--- a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryStorage.java
+++ b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryStorage.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.memory;
import org.thingsboard.server.TbQueueMsg;
diff --git a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueConsumer.java b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueConsumer.java
index a0a35eefb8..b8ddcff89c 100644
--- a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueConsumer.java
+++ b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueConsumer.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.memory;
import org.thingsboard.server.TbQueueConsumer;
diff --git a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueProducer.java b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueProducer.java
index 32b9399164..c16d34a81a 100644
--- a/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueProducer.java
+++ b/common/queue/src/main/java/org/thingsboard/server/memory/InMemoryTbQueueProducer.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.memory;
import com.google.common.util.concurrent.Futures;
diff --git a/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTbCoreQueueProvider.java
index 2f9a9ad8f8..62bc2dd5ac 100644
--- a/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTbCoreQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTbCoreQueueProvider.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.provider;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
diff --git a/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTransportQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTransportQueueProvider.java
index 945888a45a..265373793a 100644
--- a/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTransportQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/provider/KafkaTransportQueueProvider.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.provider;
import lombok.extern.slf4j.Slf4j;
diff --git a/common/queue/src/main/java/org/thingsboard/server/provider/TbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/provider/TbCoreQueueProvider.java
index b7e210f0f6..43878687ce 100644
--- a/common/queue/src/main/java/org/thingsboard/server/provider/TbCoreQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/provider/TbCoreQueueProvider.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.provider;
import org.thingsboard.server.TbQueueConsumer;
diff --git a/common/queue/src/main/java/org/thingsboard/server/provider/TransportQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/provider/TransportQueueProvider.java
index 14d8efd534..de69e1819a 100644
--- a/common/queue/src/main/java/org/thingsboard/server/provider/TransportQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/provider/TransportQueueProvider.java
@@ -1,3 +1,18 @@
+/**
+ * Copyright © 2016-2020 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.provider;
import org.thingsboard.server.TbQueueConsumer;
diff --git a/common/queue/src/main/proto/transport.proto b/common/queue/src/main/proto/transport.proto
index 892139a84a..8c564b6540 100644
--- a/common/queue/src/main/proto/transport.proto
+++ b/common/queue/src/main/proto/transport.proto
@@ -19,6 +19,16 @@ package transport;
option java_package = "org.thingsboard.server.gen.transport";
option java_outer_classname = "TransportProtos";
+/**
+ * Service Discovery Data Structures;
+ */
+message ServiceInfo {
+ string serviceId = 1;
+ repeated string serviceTypes = 2;
+ int64 tenantIdMSB = 3;
+ int64 tenantIdLSB = 4;
+}
+
/**
* Transport Service Data Structures;
*/
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 e4dc2ba1f1..cdf6b3f763 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
@@ -5,7 +5,7 @@
* 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
+ * 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,