From be84a0fc202d8ca138aed48bfd14c8c226e996c0 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 21 Dec 2020 16:42:45 +0200 Subject: [PATCH 01/12] Version set to 2.5.6-SNAPSHOT --- application/pom.xml | 2 +- common/actor/pom.xml | 2 +- common/dao-api/pom.xml | 2 +- common/data/pom.xml | 2 +- common/message/pom.xml | 2 +- common/pom.xml | 2 +- common/queue/pom.xml | 2 +- common/stats/pom.xml | 2 +- common/transport/coap/pom.xml | 2 +- common/transport/http/pom.xml | 2 +- common/transport/mqtt/pom.xml | 2 +- common/transport/pom.xml | 2 +- common/transport/transport-api/pom.xml | 2 +- common/util/pom.xml | 2 +- dao/pom.xml | 2 +- msa/black-box-tests/pom.xml | 2 +- msa/js-executor/package-lock.json | 14 +++---- msa/js-executor/package.json | 2 +- msa/js-executor/pom.xml | 2 +- msa/pom.xml | 2 +- msa/tb-node/pom.xml | 2 +- msa/tb/pom.xml | 2 +- msa/transport/coap/pom.xml | 2 +- msa/transport/http/pom.xml | 2 +- msa/transport/mqtt/pom.xml | 2 +- msa/transport/pom.xml | 2 +- msa/web-ui/package-lock.json | 2 +- msa/web-ui/package.json | 2 +- msa/web-ui/pom.xml | 2 +- netty-mqtt/pom.xml | 4 +- pom.xml | 2 +- rest-client/pom.xml | 2 +- rule-engine/pom.xml | 2 +- rule-engine/rule-engine-api/pom.xml | 2 +- rule-engine/rule-engine-components/pom.xml | 2 +- tools/pom.xml | 2 +- transport/coap/pom.xml | 2 +- transport/http/pom.xml | 2 +- transport/mqtt/pom.xml | 2 +- transport/pom.xml | 2 +- ui/package-lock.json | 43 ++++++---------------- ui/package.json | 2 +- ui/pom.xml | 2 +- 43 files changed, 61 insertions(+), 80 deletions(-) diff --git a/application/pom.xml b/application/pom.xml index 1783d9f4c5..02cdfc36e7 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard application diff --git a/common/actor/pom.xml b/common/actor/pom.xml index ee626dac6f..66a9883a14 100644 --- a/common/actor/pom.xml +++ b/common/actor/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/dao-api/pom.xml b/common/dao-api/pom.xml index 879e479c3a..38528d8511 100644 --- a/common/dao-api/pom.xml +++ b/common/dao-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/data/pom.xml b/common/data/pom.xml index 8b43e5c8b1..9b5584744f 100644 --- a/common/data/pom.xml +++ b/common/data/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/message/pom.xml b/common/message/pom.xml index 62802d141c..c956e6469c 100644 --- a/common/message/pom.xml +++ b/common/message/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/pom.xml b/common/pom.xml index e0a3865e08..9421aed1b4 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard common diff --git a/common/queue/pom.xml b/common/queue/pom.xml index b302313baf..eae5af22e1 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/stats/pom.xml b/common/stats/pom.xml index 8ccf34961c..5c9129a919 100644 --- a/common/stats/pom.xml +++ b/common/stats/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/transport/coap/pom.xml b/common/transport/coap/pom.xml index a237662b0d..3c45a7000f 100644 --- a/common/transport/coap/pom.xml +++ b/common/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/http/pom.xml b/common/transport/http/pom.xml index 715cb3d169..9b60bfb470 100644 --- a/common/transport/http/pom.xml +++ b/common/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/mqtt/pom.xml b/common/transport/mqtt/pom.xml index e3579876e4..26008b6d4a 100644 --- a/common/transport/mqtt/pom.xml +++ b/common/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/transport/pom.xml b/common/transport/pom.xml index e40f3e4009..7531c614e5 100644 --- a/common/transport/pom.xml +++ b/common/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml index 739aa2b068..e04440b23e 100644 --- a/common/transport/transport-api/pom.xml +++ b/common/transport/transport-api/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.common - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.common.transport diff --git a/common/util/pom.xml b/common/util/pom.xml index fceed3ee0e..603e365243 100644 --- a/common/util/pom.xml +++ b/common/util/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT common org.thingsboard.common diff --git a/dao/pom.xml b/dao/pom.xml index 8c1cd17bf4..39475c6329 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard dao diff --git a/msa/black-box-tests/pom.xml b/msa/black-box-tests/pom.xml index d24f38bc36..226dfa913d 100644 --- a/msa/black-box-tests/pom.xml +++ b/msa/black-box-tests/pom.xml @@ -21,7 +21,7 @@ org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/js-executor/package-lock.json b/msa/js-executor/package-lock.json index 80023b25cc..5255cd79d3 100644 --- a/msa/js-executor/package-lock.json +++ b/msa/js-executor/package-lock.json @@ -1,6 +1,6 @@ { "name": "thingsboard-js-executor", - "version": "2.5.5", + "version": "2.5.6", "lockfileVersion": 1, "requires": true, "dependencies": { @@ -1826,7 +1826,7 @@ }, "get-stream": { "version": "3.0.0", - "resolved": "https://registry.npmjs.org/get-stream/-/get-stream-3.0.0.tgz", + "resolved": "http://registry.npmjs.org/get-stream/-/get-stream-3.0.0.tgz", "integrity": "sha1-jpQ9E1jcN1VQVOy+LtsFqhdO3hQ=", "dev": true }, @@ -1936,7 +1936,7 @@ }, "got": { "version": "6.7.1", - "resolved": "https://registry.npmjs.org/got/-/got-6.7.1.tgz", + "resolved": "http://registry.npmjs.org/got/-/got-6.7.1.tgz", "integrity": "sha1-JAzQV4WpoY5WHcG0S0HHY+8ejbA=", "dev": true, "requires": { @@ -2275,7 +2275,7 @@ }, "is-obj": { "version": "1.0.1", - "resolved": "https://registry.npmjs.org/is-obj/-/is-obj-1.0.1.tgz", + "resolved": "http://registry.npmjs.org/is-obj/-/is-obj-1.0.1.tgz", "integrity": "sha1-PkcprB9f3gJc19g6iW2rn09n2w8=", "dev": true }, @@ -2917,7 +2917,7 @@ }, "path-is-absolute": { "version": "1.0.1", - "resolved": "https://registry.npmjs.org/path-is-absolute/-/path-is-absolute-1.0.1.tgz", + "resolved": "http://registry.npmjs.org/path-is-absolute/-/path-is-absolute-1.0.1.tgz", "integrity": "sha1-F0uSaHNVNP+8es5r9TpanhtcX18=", "dev": true }, @@ -3474,7 +3474,7 @@ }, "safe-regex": { "version": "1.1.0", - "resolved": "https://registry.npmjs.org/safe-regex/-/safe-regex-1.1.0.tgz", + "resolved": "http://registry.npmjs.org/safe-regex/-/safe-regex-1.1.0.tgz", "integrity": "sha1-QKNmnzsHfR6UPURinhV91IAjvy4=", "dev": true, "requires": { @@ -3832,7 +3832,7 @@ }, "strip-eof": { "version": "1.0.0", - "resolved": "https://registry.npmjs.org/strip-eof/-/strip-eof-1.0.0.tgz", + "resolved": "http://registry.npmjs.org/strip-eof/-/strip-eof-1.0.0.tgz", "integrity": "sha1-u0P/VZim6wXYm1n80SnJgzE2Br8=", "dev": true }, diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 405d582f04..d56a01bd0b 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -1,7 +1,7 @@ { "name": "thingsboard-js-executor", "private": true, - "version": "2.5.5", + "version": "2.5.6", "description": "ThingsBoard JavaScript Executor Microservice", "main": "server.js", "bin": "server.js", diff --git a/msa/js-executor/pom.xml b/msa/js-executor/pom.xml index 60ccd9d5bb..c379662359 100644 --- a/msa/js-executor/pom.xml +++ b/msa/js-executor/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/pom.xml b/msa/pom.xml index c627dca51a..a7cac82052 100644 --- a/msa/pom.xml +++ b/msa/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard msa diff --git a/msa/tb-node/pom.xml b/msa/tb-node/pom.xml index 944521f87e..2213bccd20 100644 --- a/msa/tb-node/pom.xml +++ b/msa/tb-node/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/tb/pom.xml b/msa/tb/pom.xml index 349cc36337..ef2d3eaa16 100644 --- a/msa/tb/pom.xml +++ b/msa/tb/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/transport/coap/pom.xml b/msa/transport/coap/pom.xml index 3e956736a1..919ace3232 100644 --- a/msa/transport/coap/pom.xml +++ b/msa/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/http/pom.xml b/msa/transport/http/pom.xml index bc8ccc3f1f..7d492eb851 100644 --- a/msa/transport/http/pom.xml +++ b/msa/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/mqtt/pom.xml b/msa/transport/mqtt/pom.xml index b2bc3ff8d9..dc3b23bd26 100644 --- a/msa/transport/mqtt/pom.xml +++ b/msa/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard.msa - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.msa.transport diff --git a/msa/transport/pom.xml b/msa/transport/pom.xml index d00683c492..d8064680b8 100644 --- a/msa/transport/pom.xml +++ b/msa/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/msa/web-ui/package-lock.json b/msa/web-ui/package-lock.json index b7785ef48b..adcd46ec7b 100644 --- a/msa/web-ui/package-lock.json +++ b/msa/web-ui/package-lock.json @@ -1,6 +1,6 @@ { "name": "thingsboard-web-ui", - "version": "2.5.3", + "version": "2.5.6", "lockfileVersion": 1, "requires": true, "dependencies": { diff --git a/msa/web-ui/package.json b/msa/web-ui/package.json index 2eec7556f9..3d182abfc7 100644 --- a/msa/web-ui/package.json +++ b/msa/web-ui/package.json @@ -1,7 +1,7 @@ { "name": "thingsboard-web-ui", "private": true, - "version": "2.5.5", + "version": "2.5.6", "description": "ThingsBoard Web UI Microservice", "main": "server.js", "bin": "server.js", diff --git a/msa/web-ui/pom.xml b/msa/web-ui/pom.xml index 660920835b..697363cc1b 100644 --- a/msa/web-ui/pom.xml +++ b/msa/web-ui/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT msa org.thingsboard.msa diff --git a/netty-mqtt/pom.xml b/netty-mqtt/pom.xml index 58cbcf4e6a..266a9bb235 100644 --- a/netty-mqtt/pom.xml +++ b/netty-mqtt/pom.xml @@ -19,11 +19,11 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard netty-mqtt - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT jar Netty MQTT Client diff --git a/pom.xml b/pom.xml index 6974497db2..2030a69179 100755 --- a/pom.xml +++ b/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT pom Thingsboard diff --git a/rest-client/pom.xml b/rest-client/pom.xml index 061efa390b..4dafeaabae 100644 --- a/rest-client/pom.xml +++ b/rest-client/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard rest-client diff --git a/rule-engine/pom.xml b/rule-engine/pom.xml index 00d445cc3c..98d7054f8c 100644 --- a/rule-engine/pom.xml +++ b/rule-engine/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard rule-engine diff --git a/rule-engine/rule-engine-api/pom.xml b/rule-engine/rule-engine-api/pom.xml index 2c1e323165..a9b25c4f88 100644 --- a/rule-engine/rule-engine-api/pom.xml +++ b/rule-engine/rule-engine-api/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT rule-engine org.thingsboard.rule-engine diff --git a/rule-engine/rule-engine-components/pom.xml b/rule-engine/rule-engine-components/pom.xml index 5bb14db250..23faf7707e 100644 --- a/rule-engine/rule-engine-components/pom.xml +++ b/rule-engine/rule-engine-components/pom.xml @@ -22,7 +22,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT rule-engine org.thingsboard.rule-engine diff --git a/tools/pom.xml b/tools/pom.xml index 10bfa3cc5b..763e4a98cc 100644 --- a/tools/pom.xml +++ b/tools/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard tools diff --git a/transport/coap/pom.xml b/transport/coap/pom.xml index eeea5a378a..9c53882520 100644 --- a/transport/coap/pom.xml +++ b/transport/coap/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/http/pom.xml b/transport/http/pom.xml index 7540b78507..6255717009 100644 --- a/transport/http/pom.xml +++ b/transport/http/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/mqtt/pom.xml b/transport/mqtt/pom.xml index 6f4edd87b1..b5294309df 100644 --- a/transport/mqtt/pom.xml +++ b/transport/mqtt/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT transport org.thingsboard.transport diff --git a/transport/pom.xml b/transport/pom.xml index 59b5b2470b..ca50c573d2 100644 --- a/transport/pom.xml +++ b/transport/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard transport diff --git a/ui/package-lock.json b/ui/package-lock.json index 709137d524..185fcc140f 100644 --- a/ui/package-lock.json +++ b/ui/package-lock.json @@ -1,6 +1,6 @@ { "name": "thingsboard", - "version": "2.5.3", + "version": "2.5.6", "lockfileVersion": 1, "requires": true, "dependencies": { @@ -5238,8 +5238,7 @@ "ansi-regex": { "version": "2.1.1", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "aproba": { "version": "1.2.0", @@ -5260,14 +5259,12 @@ "balanced-match": { "version": "1.0.0", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "brace-expansion": { "version": "1.1.11", "bundled": true, "dev": true, - "optional": true, "requires": { "balanced-match": "^1.0.0", "concat-map": "0.0.1" @@ -5282,20 +5279,17 @@ "code-point-at": { "version": "1.1.0", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "concat-map": { "version": "0.0.1", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "console-control-strings": { "version": "1.1.0", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "core-util-is": { "version": "1.0.2", @@ -5412,8 +5406,7 @@ "inherits": { "version": "2.0.4", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "ini": { "version": "1.3.5", @@ -5425,7 +5418,6 @@ "version": "1.0.0", "bundled": true, "dev": true, - "optional": true, "requires": { "number-is-nan": "^1.0.0" } @@ -5440,7 +5432,6 @@ "version": "3.0.4", "bundled": true, "dev": true, - "optional": true, "requires": { "brace-expansion": "^1.1.7" } @@ -5448,14 +5439,12 @@ "minimist": { "version": "0.0.8", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "minipass": { "version": "2.9.0", "bundled": true, "dev": true, - "optional": true, "requires": { "safe-buffer": "^5.1.2", "yallist": "^3.0.0" @@ -5474,7 +5463,6 @@ "version": "0.5.1", "bundled": true, "dev": true, - "optional": true, "requires": { "minimist": "0.0.8" } @@ -5564,8 +5552,7 @@ "number-is-nan": { "version": "1.0.1", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "object-assign": { "version": "4.1.1", @@ -5577,7 +5564,6 @@ "version": "1.4.0", "bundled": true, "dev": true, - "optional": true, "requires": { "wrappy": "1" } @@ -5663,8 +5649,7 @@ "safe-buffer": { "version": "5.1.2", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "safer-buffer": { "version": "2.1.2", @@ -5700,7 +5685,6 @@ "version": "1.0.2", "bundled": true, "dev": true, - "optional": true, "requires": { "code-point-at": "^1.0.0", "is-fullwidth-code-point": "^1.0.0", @@ -5720,7 +5704,6 @@ "version": "3.0.1", "bundled": true, "dev": true, - "optional": true, "requires": { "ansi-regex": "^2.0.0" } @@ -5764,14 +5747,12 @@ "wrappy": { "version": "1.0.2", "bundled": true, - "dev": true, - "optional": true + "dev": true }, "yallist": { "version": "3.1.1", "bundled": true, - "dev": true, - "optional": true + "dev": true } } }, diff --git a/ui/package.json b/ui/package.json index cbc4b8c6af..a872387289 100644 --- a/ui/package.json +++ b/ui/package.json @@ -1,7 +1,7 @@ { "name": "thingsboard", "private": true, - "version": "2.5.5", + "version": "2.5.6", "description": "ThingsBoard UI", "licenses": [ { diff --git a/ui/pom.xml b/ui/pom.xml index b27f492c97..f277bb09d8 100644 --- a/ui/pom.xml +++ b/ui/pom.xml @@ -20,7 +20,7 @@ 4.0.0 org.thingsboard - 2.5.5-SNAPSHOT + 2.5.6-SNAPSHOT thingsboard org.thingsboard From 0de5868bc5886e3c2acc2cc1c56dba8dd0fa9a7b Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 12:36:35 +0200 Subject: [PATCH 02/12] added Caffeine cache for Cassandra ts partitions saving --- .../src/main/resources/thingsboard.yml | 3 +- .../CassandraBaseTimeseriesDao.java | 46 ++++++ .../CassandraPartitionCacheKey.java | 30 ++++ .../CassandraTsPartitionsCache.java | 42 ++++++ .../nosql/CassandraPartitionsCacheTest.java | 132 ++++++++++++++++++ .../test/resources/cassandra-test.properties | 2 + 6 files changed, 254 insertions(+), 1 deletion(-) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java create mode 100644 dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 4cefe63e54..bf5ca8023d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -223,8 +223,9 @@ cassandra: read_consistency_level: "${CASSANDRA_READ_CONSISTENCY_LEVEL:ONE}" write_consistency_level: "${CASSANDRA_WRITE_CONSISTENCY_LEVEL:ONE}" default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}" - # Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS,INDEFINITE + # Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS, INDEFINITE ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" + ts_key_value_partitions_max_cache_size: "${TS_KV_PARTITIONS_MAX_CACHE_SIZE:100000}" ts_key_value_ttl: "${TS_KV_TTL:0}" events_ttl: "${TS_EVENTS_TTL:0}" # Specify TTL of debug log in seconds. The current value corresponds to one week diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index b96462e350..d39c6c5afb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -29,6 +29,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -66,6 +67,7 @@ import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -88,12 +90,17 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem public static final String DESC_ORDER = "DESC"; private static List FIXED_PARTITION = Arrays.asList(new Long[]{0L}); + private CassandraTsPartitionsCache cassandraTsPartitionsCache; + @Autowired private Environment environment; @Value("${cassandra.query.ts_key_value_partitioning}") private String partitioning; + @Value("${cassandra.query.ts_key_value_partitions_max_cache_size}") + private long partitionsCacheSize; + @Value("${cassandra.query.ts_key_value_ttl}") private long systemTtl; @@ -126,6 +133,9 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem Optional partition = NoSqlTsPartitionDate.parse(partitioning); if (partition.isPresent()) { tsFormat = partition.get(); + if (!isFixedPartitioning() && partitionsCacheSize > 0) { + cassandraTsPartitionsCache = new CassandraTsPartitionsCache(partitionsCacheSize); + } } else { log.warn("Incorrect configuration of partitioning {}", partitioning); throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); @@ -390,6 +400,42 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } ttl = computeTtl(ttl); long partition = toPartitionTs(tsKvEntryTs); + if (cassandraTsPartitionsCache == null) { + return doSavePartition(tenantId, entityId, key, ttl, partition); + } else { + CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); + CompletableFuture hasFuture = cassandraTsPartitionsCache.has(partitionSearchKey); + SettableFuture listenableFuture = SettableFuture.create(); + if (hasFuture == null) { + return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, ttl); + } else { + hasFuture.whenComplete((result, throwable) -> { + if (throwable != null) { + listenableFuture.setException(throwable); + } else { + listenableFuture.set(result); + } + }); + long finalTtl = ttl; + return Futures.transformAsync(listenableFuture, result -> { + if (result) { + return Futures.immediateFuture(null); + } else { + return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, finalTtl); + } + }, readResultsProcessingExecutor); + } + } + } + + private ListenableFuture processDoSavePartition(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { + return Futures.transformAsync(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { + cassandraTsPartitionsCache.put(partitionSearchKey); + return Futures.immediateFuture(input); + }, readResultsProcessingExecutor); + } + + private ListenableFuture doSavePartition(TenantId tenantId, EntityId entityId, String key, long ttl, long partition) { log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); BoundStatement stmt = (ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind(); stmt = stmt.setString(0, entityId.getEntityType().name()) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java new file mode 100644 index 0000000000..791ce84113 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.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.dao.timeseries; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.id.EntityId; + +@Data +@AllArgsConstructor +public class CassandraPartitionCacheKey { + + private EntityId entityId; + private String key; + private long partition; + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java new file mode 100644 index 0000000000..b467b5446f --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java @@ -0,0 +1,42 @@ +/** + * 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.dao.timeseries; + +import com.github.benmanes.caffeine.cache.AsyncLoadingCache; +import com.github.benmanes.caffeine.cache.Caffeine; + +import java.util.concurrent.CompletableFuture; + +public class CassandraTsPartitionsCache { + + private AsyncLoadingCache partitionsCache; + + public CassandraTsPartitionsCache(long maxCacheSize) { + this.partitionsCache = Caffeine.newBuilder() + .maximumSize(maxCacheSize) + .buildAsync(key -> { + throw new IllegalStateException("'get' methods calls are not supported!"); + }); + } + + public CompletableFuture has(CassandraPartitionCacheKey key) { + return partitionsCache.getIfPresent(key); + } + + public void put(CassandraPartitionCacheKey key) { + partitionsCache.put(key, CompletableFuture.completedFuture(true)); + } +} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java new file mode 100644 index 0000000000..82c0e0d6d2 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -0,0 +1,132 @@ +/** + * 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.dao.nosql; + +import com.datastax.driver.core.BoundStatement; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.CodecRegistry; +import com.datastax.driver.core.Configuration; +import com.datastax.driver.core.ConsistencyLevel; +import com.datastax.driver.core.PreparedStatement; +import com.datastax.driver.core.ResultSetFuture; +import com.datastax.driver.core.Session; +import com.datastax.driver.core.Statement; +import com.google.common.util.concurrent.Futures; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.core.env.Environment; +import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.cassandra.CassandraCluster; +import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao; + +import java.util.UUID; + +import static org.mockito.Matchers.any; +import static org.mockito.Matchers.anyInt; +import static org.mockito.Matchers.anyLong; +import static org.mockito.Matchers.anyString; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.class) +public class CassandraPartitionsCacheTest { + + private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao; + + @Mock + private Environment environment; + + @Mock + private CassandraBufferedRateExecutor rateLimiter; + + @Mock + private CassandraCluster cluster; + + @Mock + private Session session; + + @Mock + private Cluster sessionCluster; + + @Mock + private Configuration configuration; + + @Mock + private PreparedStatement preparedStatement; + + @Mock + private BoundStatement boundStatement; + + @Before + public void setUp() { + when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getSession()).thenReturn(session); + when(session.getCluster()).thenReturn(sessionCluster); + when(sessionCluster.getConfiguration()).thenReturn(configuration); + when(configuration.getCodecRegistry()).thenReturn(CodecRegistry.DEFAULT_INSTANCE); + when(session.prepare(anyString())).thenReturn(preparedStatement); + when(preparedStatement.bind()).thenReturn(boundStatement); + when(boundStatement.setString(anyInt(), anyString())).thenReturn(boundStatement); + when(boundStatement.setUUID(anyInt(), any(UUID.class))).thenReturn(boundStatement); + when(boundStatement.setLong(anyInt(), anyLong())).thenReturn(boundStatement); + when(boundStatement.setInt(anyInt(), anyInt())).thenReturn(boundStatement); + + cassandraBaseTimeseriesDao = spy(new CassandraBaseTimeseriesDao()); + + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitioning", "MONTHS"); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitionsCacheSize", 100000); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "systemTtl", 0); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "setNullValuesEnabled", false); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "environment", environment); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "rateLimiter", rateLimiter); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "cluster", cluster); + + doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(ResultSetFuture.class), any()); + + } + + @Test + public void testPartitionSave() throws Exception { + + cassandraBaseTimeseriesDao.init(); + + + UUID id = UUID.randomUUID(); + TenantId tenantId = new TenantId(id); + long tsKvEntryTs = System.currentTimeMillis(); + + for (int i = 0; i < 50000; i++) { + cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); + } + + for (int i = 0; i < 60000; i++) { + cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); + } + + verify(cassandraBaseTimeseriesDao, times(60000)).executeAsyncWrite(any(TenantId.class), any(Statement.class)); + + + } + +} \ No newline at end of file diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 51f34a08d6..2bf57b190f 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -46,6 +46,8 @@ cassandra.query.default_fetch_size=2000 cassandra.query.ts_key_value_partitioning=HOURS +cassandra.query.ts_key_value_partitions_max_cache_size=100000 + cassandra.query.ts_key_value_ttl=0 cassandra.query.debug_events_ttl=604800 From aede1af6f947444144b3b1f64a0681bccbdab57e Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 14:50:07 +0200 Subject: [PATCH 03/12] improvements after review --- .../CassandraBaseTimeseriesDao.java | 20 ++++++++----------- .../nosql/CassandraPartitionsCacheTest.java | 7 ------- 2 files changed, 8 insertions(+), 19 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index d39c6c5afb..7d834a8eea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -405,33 +405,29 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } else { CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); CompletableFuture hasFuture = cassandraTsPartitionsCache.has(partitionSearchKey); - SettableFuture listenableFuture = SettableFuture.create(); + SettableFuture listenableFuture = SettableFuture.create(); if (hasFuture == null) { return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, ttl); } else { + long finalTtl = ttl; hasFuture.whenComplete((result, throwable) -> { if (throwable != null) { listenableFuture.setException(throwable); + } else if (result) { + listenableFuture.set(null); } else { - listenableFuture.set(result); + listenableFuture.setFuture(processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, finalTtl)); } }); - long finalTtl = ttl; - return Futures.transformAsync(listenableFuture, result -> { - if (result) { - return Futures.immediateFuture(null); - } else { - return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, finalTtl); - } - }, readResultsProcessingExecutor); + return listenableFuture; } } } private ListenableFuture processDoSavePartition(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { - return Futures.transformAsync(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { + return Futures.transform(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { cassandraTsPartitionsCache.put(partitionSearchKey); - return Futures.immediateFuture(input); + return input; }, readResultsProcessingExecutor); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java index 82c0e0d6d2..a67d84964a 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -103,15 +103,12 @@ public class CassandraPartitionsCacheTest { ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "cluster", cluster); doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(ResultSetFuture.class), any()); - } @Test public void testPartitionSave() throws Exception { - cassandraBaseTimeseriesDao.init(); - UUID id = UUID.randomUUID(); TenantId tenantId = new TenantId(id); long tsKvEntryTs = System.currentTimeMillis(); @@ -119,14 +116,10 @@ public class CassandraPartitionsCacheTest { for (int i = 0; i < 50000; i++) { cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); } - for (int i = 0; i < 60000; i++) { cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); } verify(cassandraBaseTimeseriesDao, times(60000)).executeAsyncWrite(any(TenantId.class), any(Statement.class)); - - } - } \ No newline at end of file From 724df48f7911561ae1c6dc41d4f7fbe0c81d9be1 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 15:06:15 +0200 Subject: [PATCH 04/12] added default value for .yml parameter & rename params names --- .../CassandraBaseTimeseriesDao.java | 22 +++++++++---------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 7d834a8eea..a9671cad6e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -98,7 +98,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem @Value("${cassandra.query.ts_key_value_partitioning}") private String partitioning; - @Value("${cassandra.query.ts_key_value_partitions_max_cache_size}") + @Value("${cassandra.query.ts_key_value_partitions_max_cache_size:100000}") private long partitionsCacheSize; @Value("${cassandra.query.ts_key_value_ttl}") @@ -404,27 +404,27 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return doSavePartition(tenantId, entityId, key, ttl, partition); } else { CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); - CompletableFuture hasFuture = cassandraTsPartitionsCache.has(partitionSearchKey); - SettableFuture listenableFuture = SettableFuture.create(); - if (hasFuture == null) { - return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, ttl); + CompletableFuture hasInCacheFuture = cassandraTsPartitionsCache.has(partitionSearchKey); + SettableFuture futureResult = SettableFuture.create(); + if (hasInCacheFuture == null) { + return doSavePartitionWithCache(tenantId, entityId, key, partition, partitionSearchKey, ttl); } else { long finalTtl = ttl; - hasFuture.whenComplete((result, throwable) -> { + hasInCacheFuture.whenComplete((result, throwable) -> { if (throwable != null) { - listenableFuture.setException(throwable); + futureResult.setException(throwable); } else if (result) { - listenableFuture.set(null); + futureResult.set(null); } else { - listenableFuture.setFuture(processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, finalTtl)); + futureResult.setFuture(doSavePartitionWithCache(tenantId, entityId, key, partition, partitionSearchKey, finalTtl)); } }); - return listenableFuture; + return futureResult; } } } - private ListenableFuture processDoSavePartition(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { + private ListenableFuture doSavePartitionWithCache(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { return Futures.transform(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { cassandraTsPartitionsCache.put(partitionSearchKey); return input; From 0a2255d0557f8881bc0f6d95bb6a085762fa2989 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 15:46:16 +0200 Subject: [PATCH 05/12] move SettableFuture.create() to the else block --- .../server/dao/timeseries/CassandraBaseTimeseriesDao.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index a9671cad6e..46f8b4a5cb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -405,11 +405,11 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } else { CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); CompletableFuture hasInCacheFuture = cassandraTsPartitionsCache.has(partitionSearchKey); - SettableFuture futureResult = SettableFuture.create(); if (hasInCacheFuture == null) { return doSavePartitionWithCache(tenantId, entityId, key, partition, partitionSearchKey, ttl); } else { long finalTtl = ttl; + SettableFuture futureResult = SettableFuture.create(); hasInCacheFuture.whenComplete((result, throwable) -> { if (throwable != null) { futureResult.setException(throwable); From 2ea3b18738e95e646b6220fb898781d53d9f2e8b Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 21 Dec 2020 17:05:12 +0200 Subject: [PATCH 06/12] partitions cache improvements --- .../CassandraBaseTimeseriesDao.java | 41 ++++++++++--------- .../CassandraTsPartitionsCache.java | 4 +- 2 files changed, 23 insertions(+), 22 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 46f8b4a5cb..7bcae3b740 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -404,31 +404,32 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return doSavePartition(tenantId, entityId, key, ttl, partition); } else { CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); - CompletableFuture hasInCacheFuture = cassandraTsPartitionsCache.has(partitionSearchKey); - if (hasInCacheFuture == null) { - return doSavePartitionWithCache(tenantId, entityId, key, partition, partitionSearchKey, ttl); + if (!cassandraTsPartitionsCache.has(partitionSearchKey)) { + ListenableFuture result = doSavePartition(tenantId, entityId, key, ttl, partition); + Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor()); + return result; } else { - long finalTtl = ttl; - SettableFuture futureResult = SettableFuture.create(); - hasInCacheFuture.whenComplete((result, throwable) -> { - if (throwable != null) { - futureResult.setException(throwable); - } else if (result) { - futureResult.set(null); - } else { - futureResult.setFuture(doSavePartitionWithCache(tenantId, entityId, key, partition, partitionSearchKey, finalTtl)); - } - }); - return futureResult; + return Futures.immediateFuture(null); } } } - private ListenableFuture doSavePartitionWithCache(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { - return Futures.transform(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { - cassandraTsPartitionsCache.put(partitionSearchKey); - return input; - }, readResultsProcessingExecutor); + private class CacheCallback implements FutureCallback { + private final CassandraPartitionCacheKey key; + + private CacheCallback(CassandraPartitionCacheKey key) { + this.key = key; + } + + @Override + public void onSuccess(Void result) { + cassandraTsPartitionsCache.put(key); + } + + @Override + public void onFailure(Throwable t) { + + } } private ListenableFuture doSavePartition(TenantId tenantId, EntityId entityId, String key, long ttl, long partition) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java index b467b5446f..bafc00c872 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java @@ -32,8 +32,8 @@ public class CassandraTsPartitionsCache { }); } - public CompletableFuture has(CassandraPartitionCacheKey key) { - return partitionsCache.getIfPresent(key); + public boolean has(CassandraPartitionCacheKey key) { + return partitionsCache.getIfPresent(key) != null; } public void put(CassandraPartitionCacheKey key) { From bff16ddafd650f5f1117fd0a73131ee5e07bc133 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 22 Dec 2020 18:02:38 +0200 Subject: [PATCH 07/12] Mqtt flag sessionPresent depends on flag isCleanSession --- .../transport/mqtt/MqttTransportHandler.java | 50 +++++++++---------- 1 file changed, 25 insertions(+), 25 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index abb8f629f1..7f45577536 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -357,7 +357,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier()); X509Certificate cert; if (sslHandler != null && (cert = getX509Certificate()) != null) { - processX509CertConnect(ctx, cert); + processX509CertConnect(ctx, cert, msg); } else { processAuthTokenConnect(ctx, msg); } @@ -367,27 +367,27 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement String userName = msg.payload().userName(); log.info("[{}] Processing connect msg for client with user name: {}!", sessionId, userName); if (StringUtils.isEmpty(userName)) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD, msg)); ctx.close(); } else { transportService.process(ValidateDeviceTokenRequestMsg.newBuilder().setToken(userName).build(), new TransportServiceCallback() { @Override - public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { - onValidateDeviceResponse(msg, ctx); + public void onSuccess(ValidateDeviceCredentialsResponseMsg responseMsg) { + onValidateDeviceResponse(responseMsg, ctx, msg); } @Override public void onError(Throwable e) { log.trace("[{}] Failed to process credentials: {}", address, userName, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); ctx.close(); } }); } } - private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert) { + private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert, MqttConnectMessage msg) { try { if(!context.isSkipValidityCheckForClientCert()){ cert.checkValidity(); @@ -397,19 +397,19 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(ValidateDeviceX509CertRequestMsg.newBuilder().setHash(sha3Hash).build(), new TransportServiceCallback() { @Override - public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { - onValidateDeviceResponse(msg, ctx); + public void onSuccess(ValidateDeviceCredentialsResponseMsg responseMsg) { + onValidateDeviceResponse(responseMsg, ctx, msg); } @Override public void onError(Throwable e) { log.trace("[{}] Failed to process credentials: {}", address, sha3Hash, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); ctx.close(); } }); } catch (Exception e) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, msg)); ctx.close(); } } @@ -433,11 +433,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement doDisconnect(); } - private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode) { + private MqttConnAckMessage createMqttConnAckMsg(MqttConnectReturnCode returnCode, MqttConnectMessage msg) { MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(CONNACK, false, AT_MOST_ONCE, false, 0); MqttConnAckVariableHeader mqttConnAckVariableHeader = - new MqttConnAckVariableHeader(returnCode, true); + new MqttConnAckVariableHeader(returnCode, !msg.variableHeader().isCleanSession()); return new MqttConnAckMessage(mqttFixedHeader, mqttConnAckVariableHeader); } @@ -513,36 +513,36 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void onValidateDeviceResponse(ValidateDeviceCredentialsResponseMsg msg, ChannelHandlerContext ctx) { - if (!msg.hasDeviceInfo()) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED)); + private void onValidateDeviceResponse(ValidateDeviceCredentialsResponseMsg responseMsg, ChannelHandlerContext ctx, MqttConnectMessage msg) { + if (!responseMsg.hasDeviceInfo()) { + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, msg)); ctx.close(); } else { - deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo()); + deviceSessionCtx.setDeviceInfo(responseMsg.getDeviceInfo()); sessionInfo = SessionInfoProto.newBuilder() .setNodeId(context.getNodeId()) .setSessionIdMSB(sessionId.getMostSignificantBits()) .setSessionIdLSB(sessionId.getLeastSignificantBits()) - .setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB()) - .setDeviceIdLSB(msg.getDeviceInfo().getDeviceIdLSB()) - .setTenantIdMSB(msg.getDeviceInfo().getTenantIdMSB()) - .setTenantIdLSB(msg.getDeviceInfo().getTenantIdLSB()) - .setDeviceName(msg.getDeviceInfo().getDeviceName()) - .setDeviceType(msg.getDeviceInfo().getDeviceType()) + .setDeviceIdMSB(responseMsg.getDeviceInfo().getDeviceIdMSB()) + .setDeviceIdLSB(responseMsg.getDeviceInfo().getDeviceIdLSB()) + .setTenantIdMSB(responseMsg.getDeviceInfo().getTenantIdMSB()) + .setTenantIdLSB(responseMsg.getDeviceInfo().getTenantIdLSB()) + .setDeviceName(responseMsg.getDeviceInfo().getDeviceName()) + .setDeviceType(responseMsg.getDeviceInfo().getDeviceType()) .build(); transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), new TransportServiceCallback() { @Override - public void onSuccess(Void msg) { + public void onSuccess(Void response) { transportService.registerAsyncSession(sessionInfo, MqttTransportHandler.this); checkGatewaySession(); - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, msg)); log.info("[{}] Client connected!", sessionId); } @Override public void onError(Throwable e) { log.warn("[{}] Failed to submit session event", sessionId, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); ctx.close(); } }); From 16d82e77fc5000a19d30797ef195ac30b2cf406a Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 23 Dec 2020 16:54:14 +0200 Subject: [PATCH 08/12] refactored --- .../transport/mqtt/MqttTransportHandler.java | 46 +++++++++---------- 1 file changed, 23 insertions(+), 23 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 7f45577536..bfd32eb751 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -363,31 +363,31 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) { - String userName = msg.payload().userName(); + private void processAuthTokenConnect(ChannelHandlerContext ctx, MqttConnectMessage connectMessage) { + String userName = connectMessage.payload().userName(); log.info("[{}] Processing connect msg for client with user name: {}!", sessionId, userName); if (StringUtils.isEmpty(userName)) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_BAD_USER_NAME_OR_PASSWORD, connectMessage)); ctx.close(); } else { transportService.process(ValidateDeviceTokenRequestMsg.newBuilder().setToken(userName).build(), new TransportServiceCallback() { @Override - public void onSuccess(ValidateDeviceCredentialsResponseMsg responseMsg) { - onValidateDeviceResponse(responseMsg, ctx, msg); + public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { + onValidateDeviceResponse(msg, ctx, connectMessage); } @Override public void onError(Throwable e) { log.trace("[{}] Failed to process credentials: {}", address, userName, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage)); ctx.close(); } }); } } - private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert, MqttConnectMessage msg) { + private void processX509CertConnect(ChannelHandlerContext ctx, X509Certificate cert, MqttConnectMessage connectMessage) { try { if(!context.isSkipValidityCheckForClientCert()){ cert.checkValidity(); @@ -397,19 +397,19 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(ValidateDeviceX509CertRequestMsg.newBuilder().setHash(sha3Hash).build(), new TransportServiceCallback() { @Override - public void onSuccess(ValidateDeviceCredentialsResponseMsg responseMsg) { - onValidateDeviceResponse(responseMsg, ctx, msg); + public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) { + onValidateDeviceResponse(msg, ctx, connectMessage); } @Override public void onError(Throwable e) { log.trace("[{}] Failed to process credentials: {}", address, sha3Hash, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage)); ctx.close(); } }); } catch (Exception e) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage)); ctx.close(); } } @@ -513,36 +513,36 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void onValidateDeviceResponse(ValidateDeviceCredentialsResponseMsg responseMsg, ChannelHandlerContext ctx, MqttConnectMessage msg) { - if (!responseMsg.hasDeviceInfo()) { - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, msg)); + private void onValidateDeviceResponse(ValidateDeviceCredentialsResponseMsg msg, ChannelHandlerContext ctx, MqttConnectMessage connectMessage) { + if (!msg.hasDeviceInfo()) { + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_REFUSED_NOT_AUTHORIZED, connectMessage)); ctx.close(); } else { - deviceSessionCtx.setDeviceInfo(responseMsg.getDeviceInfo()); + deviceSessionCtx.setDeviceInfo(msg.getDeviceInfo()); sessionInfo = SessionInfoProto.newBuilder() .setNodeId(context.getNodeId()) .setSessionIdMSB(sessionId.getMostSignificantBits()) .setSessionIdLSB(sessionId.getLeastSignificantBits()) - .setDeviceIdMSB(responseMsg.getDeviceInfo().getDeviceIdMSB()) - .setDeviceIdLSB(responseMsg.getDeviceInfo().getDeviceIdLSB()) - .setTenantIdMSB(responseMsg.getDeviceInfo().getTenantIdMSB()) - .setTenantIdLSB(responseMsg.getDeviceInfo().getTenantIdLSB()) - .setDeviceName(responseMsg.getDeviceInfo().getDeviceName()) - .setDeviceType(responseMsg.getDeviceInfo().getDeviceType()) + .setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB()) + .setDeviceIdLSB(msg.getDeviceInfo().getDeviceIdLSB()) + .setTenantIdMSB(msg.getDeviceInfo().getTenantIdMSB()) + .setTenantIdLSB(msg.getDeviceInfo().getTenantIdLSB()) + .setDeviceName(msg.getDeviceInfo().getDeviceName()) + .setDeviceType(msg.getDeviceInfo().getDeviceType()) .build(); transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), new TransportServiceCallback() { @Override public void onSuccess(Void response) { transportService.registerAsyncSession(sessionInfo, MqttTransportHandler.this); checkGatewaySession(); - ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, connectMessage)); log.info("[{}] Client connected!", sessionId); } @Override public void onError(Throwable e) { log.warn("[{}] Failed to submit session event", sessionId, e); - ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, msg)); + ctx.writeAndFlush(createMqttConnAckMsg(MqttConnectReturnCode.CONNECTION_REFUSED_SERVER_UNAVAILABLE, connectMessage)); ctx.close(); } }); From 1d09fa2bd8509760bd7fef5cf655db99598b01a0 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 29 Dec 2020 20:15:46 +0200 Subject: [PATCH 09/12] removed ServiceId from kafka consumer groupId --- .../queue/provider/KafkaMonolithQueueFactory.java | 10 +++++----- .../server/queue/provider/KafkaTbCoreQueueFactory.java | 6 +++--- .../queue/provider/KafkaTbRuleEngineQueueFactory.java | 6 +++--- .../queue/provider/KafkaTbTransportQueueFactory.java | 4 ++-- 4 files changed, 13 insertions(+), 13 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index 31ef4d3e45..6ed49d21fc 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -152,7 +152,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(ruleEngineSettings.getTopic()); consumerBuilder.clientId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("re-" + queueName + "-consumer"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(ruleEngineAdmin); return consumerBuilder.build(); @@ -164,7 +164,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("monolith-rule-engine-notifications-consumer"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -176,7 +176,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(coreSettings.getTopic()); consumerBuilder.clientId("monolith-core-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("monolith-core-consumer-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("monolith-core-consumer"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(coreAdmin); return consumerBuilder.build(); @@ -188,7 +188,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("monolith-core-notifications-consumer"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -229,7 +229,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); + responseBuilder.groupId("rule-engine-node"); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java index bf2ec6751e..7ca2f9196d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java @@ -146,7 +146,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(coreSettings.getTopic()); consumerBuilder.clientId("tb-core-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("tb-core-node-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("tb-core-node"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(coreAdmin); return consumerBuilder.build(); @@ -158,7 +158,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("tb-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("tb-core-notifications-node-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("tb-core-notifications-node"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -199,7 +199,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); + responseBuilder.groupId("rule-engine-node"); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index dc44caf3f0..3bab6d1866 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -141,7 +141,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(ruleEngineSettings.getTopic()); consumerBuilder.clientId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("re-" + queueName + "-consumer"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(ruleEngineAdmin); return consumerBuilder.build(); @@ -153,7 +153,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("tb-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("tb-rule-engine-notifications-node-" + serviceInfoProvider.getServiceId()); + consumerBuilder.groupId("tb-rule-engine-notifications-node"); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -172,7 +172,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); + responseBuilder.groupId("rule-engine-node"); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java index 505aa203d7..58623a8f0b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java @@ -91,7 +91,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(transportApiSettings.getResponsesTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("transport-api-response-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); + responseBuilder.groupId("transport-node"); responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiResponseMsg.parseFrom(msg.getData()), msg.getHeaders())); responseBuilder.admin(transportApiAdmin); @@ -132,7 +132,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("transport-api-notifications-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); + responseBuilder.groupId("transport-node"); responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToTransportMsg.parseFrom(msg.getData()), msg.getHeaders())); responseBuilder.admin(notificationAdmin); return responseBuilder.build(); From 36868537517d5b1c0d5cba633c206710e56b2475 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 30 Dec 2020 15:07:14 +0200 Subject: [PATCH 10/12] refactored --- .../server/queue/provider/KafkaMonolithQueueFactory.java | 6 +++--- .../server/queue/provider/KafkaTbCoreQueueFactory.java | 4 ++-- .../queue/provider/KafkaTbRuleEngineQueueFactory.java | 4 ++-- .../server/queue/provider/KafkaTbTransportQueueFactory.java | 4 ++-- 4 files changed, 9 insertions(+), 9 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index 6ed49d21fc..8e2ab69915 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -164,7 +164,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("monolith-rule-engine-notifications-consumer"); + consumerBuilder.groupId("monolith-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -188,7 +188,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("monolith-core-notifications-consumer"); + consumerBuilder.groupId("monolith-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -229,7 +229,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node"); + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java index 7ca2f9196d..c6acc0fd6a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java @@ -158,7 +158,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("tb-core-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("tb-core-notifications-node"); + consumerBuilder.groupId("tb-core-notifications-node-" + serviceInfoProvider.getServiceId()); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -199,7 +199,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node"); + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index 3bab6d1866..ad8d84b4b6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -153,7 +153,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName()); consumerBuilder.clientId("tb-rule-engine-notifications-consumer-" + serviceInfoProvider.getServiceId()); - consumerBuilder.groupId("tb-rule-engine-notifications-node"); + consumerBuilder.groupId("tb-rule-engine-notifications-node-" + serviceInfoProvider.getServiceId()); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(notificationAdmin); return consumerBuilder.build(); @@ -172,7 +172,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("js-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("rule-engine-node"); + responseBuilder.groupId("rule-engine-node-" + serviceInfoProvider.getServiceId()); responseBuilder.decoder(msg -> { JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder(); JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java index 58623a8f0b..505aa203d7 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbTransportQueueFactory.java @@ -91,7 +91,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(transportApiSettings.getResponsesTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("transport-api-response-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("transport-node"); + responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiResponseMsg.parseFrom(msg.getData()), msg.getHeaders())); responseBuilder.admin(transportApiAdmin); @@ -132,7 +132,7 @@ public class KafkaTbTransportQueueFactory implements TbTransportQueueFactory { responseBuilder.settings(kafkaSettings); responseBuilder.topic(transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId()); responseBuilder.clientId("transport-api-notifications-" + serviceInfoProvider.getServiceId()); - responseBuilder.groupId("transport-node"); + responseBuilder.groupId("transport-node-" + serviceInfoProvider.getServiceId()); responseBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToTransportMsg.parseFrom(msg.getData()), msg.getHeaders())); responseBuilder.admin(notificationAdmin); return responseBuilder.build(); From 143e9cb877fb6f6e5b885184643b47d2c48e5905 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Thu, 14 Jan 2021 10:21:51 +0200 Subject: [PATCH 11/12] added new methods to JacksonUtil & replaced objectMapper from AuditLogServiceImpl --- .../server/dao/audit/AuditLogServiceImpl.java | 16 +++++++--------- .../server/dao/util/mapping/JacksonUtil.java | 9 +++++++++ 2 files changed, 16 insertions(+), 9 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java index 0d1c708f68..6c8e2a25ff 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java @@ -17,7 +17,6 @@ package org.thingsboard.server.dao.audit; import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.collect.Lists; @@ -48,6 +47,7 @@ import org.thingsboard.server.dao.audit.sink.AuditLogSink; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; +import org.thingsboard.server.dao.util.mapping.JacksonUtil; import java.io.PrintWriter; import java.io.StringWriter; @@ -61,8 +61,6 @@ import static org.thingsboard.server.dao.service.Validator.validateId; @ConditionalOnProperty(prefix = "audit-log", value = "enabled", havingValue = "true") public class AuditLogServiceImpl implements AuditLogService { - private static final ObjectMapper objectMapper = new ObjectMapper(); - private static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; private static final int INSERTS_PER_ENTRY = 3; @@ -159,7 +157,7 @@ public class AuditLogServiceImpl implements AuditLogService { private JsonNode constructActionData(I entityId, E entity, ActionType actionType, Object... additionalInfo) { - ObjectNode actionData = objectMapper.createObjectNode(); + ObjectNode actionData = JacksonUtil.newObjectNode(); switch (actionType) { case ADDED: case UPDATED: @@ -168,7 +166,7 @@ public class AuditLogServiceImpl implements AuditLogService { case RELATIONS_DELETED: case ASSIGNED_TO_TENANT: if (entity != null) { - ObjectNode entityNode = objectMapper.valueToTree(entity); + ObjectNode entityNode = (ObjectNode) JacksonUtil.valueToTree(entity); if (entityId.getEntityType() == EntityType.DASHBOARD) { entityNode.put("configuration", ""); } @@ -177,7 +175,7 @@ public class AuditLogServiceImpl implements AuditLogService { if (entityId.getEntityType() == EntityType.RULE_CHAIN) { RuleChainMetaData ruleChainMetaData = extractParameter(RuleChainMetaData.class, additionalInfo); if (ruleChainMetaData != null) { - ObjectNode ruleChainMetaDataNode = objectMapper.valueToTree(ruleChainMetaData); + ObjectNode ruleChainMetaDataNode = (ObjectNode) JacksonUtil.valueToTree(ruleChainMetaData); actionData.set("metadata", ruleChainMetaDataNode); } } @@ -194,7 +192,7 @@ public class AuditLogServiceImpl implements AuditLogService { String scope = extractParameter(String.class, 0, additionalInfo); List attributes = extractParameter(List.class, 1, additionalInfo); actionData.put("scope", scope); - ObjectNode attrsNode = objectMapper.createObjectNode(); + ObjectNode attrsNode = JacksonUtil.newObjectNode(); if (attributes != null) { for (AttributeKvEntry attr : attributes) { attrsNode.put(attr.getKey(), attr.getValueAsString()); @@ -225,7 +223,7 @@ public class AuditLogServiceImpl implements AuditLogService { case CREDENTIALS_UPDATED: actionData.put("entityId", entityId.toString()); DeviceCredentials deviceCredentials = extractParameter(DeviceCredentials.class, additionalInfo); - actionData.set("credentials", objectMapper.valueToTree(deviceCredentials)); + actionData.set("credentials", JacksonUtil.valueToTree(deviceCredentials)); break; case ASSIGNED_TO_CUSTOMER: strEntityId = extractParameter(String.class, 0, additionalInfo); @@ -246,7 +244,7 @@ public class AuditLogServiceImpl implements AuditLogService { case RELATION_ADD_OR_UPDATE: case RELATION_DELETED: EntityRelation relation = extractParameter(EntityRelation.class, 0, additionalInfo); - actionData.set("relation", objectMapper.valueToTree(relation)); + actionData.set("relation", JacksonUtil.valueToTree(relation)); break; case LOGIN: case LOGOUT: diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java index d654e85d5c..fd40d71189 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.util.mapping; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; import java.io.IOException; @@ -66,7 +67,15 @@ public class JacksonUtil { } } + public static ObjectNode newObjectNode(){ + return OBJECT_MAPPER.createObjectNode(); + } + public static T clone(T value) { return fromString(toString(value), (Class) value.getClass()); } + + public static JsonNode valueToTree(T value) { + return OBJECT_MAPPER.valueToTree(value); + } } From 6682d66be3397d74dd64ac6aed0d660072a16f61 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Sun, 17 Jan 2021 00:47:02 +0200 Subject: [PATCH 12/12] Make device private, if edge is public --- .../server/controller/EdgeController.java | 3 +++ .../edge/rpc/processor/DeviceProcessor.java | 18 ++++++++++++++++-- 2 files changed, 19 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java index e31e711e5c..15d905bfed 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java @@ -254,6 +254,9 @@ public class EdgeController extends BaseController { Customer publicCustomer = customerService.findOrCreatePublicCustomer(edge.getTenantId()); Edge savedEdge = checkNotNull(edgeService.assignEdgeToCustomer(getCurrentUser().getTenantId(), edgeId, publicCustomer.getId())); + tbClusterService.onEntityStateChange(getTenantId(), edgeId, + ComponentLifecycleEvent.UPDATED); + logEntityAction(edgeId, savedEdge, savedEdge.getCustomerId(), ActionType.ASSIGNED_TO_CUSTOMER, null, strEdgeId, publicCustomer.getId().toString(), publicCustomer.getName()); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java index 9aab6ec124..8d536d5076 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java @@ -25,9 +25,9 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang.RandomStringUtils; import org.apache.commons.lang.StringUtils; import org.checkerframework.checker.nullness.qual.Nullable; -import org.springframework.security.core.parameters.P; import org.springframework.stereotype.Component; import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.edge.Edge; @@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.util.mapping.JacksonUtil; import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.DeviceRpcCallMsg; @@ -179,7 +180,8 @@ public class DeviceProcessor extends BaseProcessor { log.debug("[{}] Creating device entity [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); device = new Device(); device.setTenantId(edge.getTenantId()); - device.setCustomerId(edge.getCustomerId()); + // make device private, if edge is public + device.setCustomerId(getCustomerId(edge)); device.setName(deviceName); device.setType(deviceUpdateMsg.getType()); device.setLabel(deviceUpdateMsg.getLabel()); @@ -194,6 +196,18 @@ public class DeviceProcessor extends BaseProcessor { return device; } + private CustomerId getCustomerId(Edge edge) { + if (edge.getCustomerId() == null || edge.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { + return edge.getCustomerId(); + } + Customer publicCustomer = customerService.findOrCreatePublicCustomer(edge.getTenantId()); + if (publicCustomer.getId().equals(edge.getCustomerId())) { + return null; + } else { + return edge.getCustomerId(); + } + } + private void createRelationFromEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId) { EntityRelation relation = new EntityRelation(); relation.setFrom(edgeId);