2144 changed files with 51712 additions and 167049 deletions
@ -0,0 +1,60 @@ |
|||
--- |
|||
name: "\U0001F41E Bug report" |
|||
about: Create a report to help us improve |
|||
title: "[Bug] " |
|||
labels: bug |
|||
assignees: ashvayka, vvlladd28 |
|||
|
|||
--- |
|||
|
|||
**Describe the bug** |
|||
A clear and concise description of what the bug is. |
|||
|
|||
**Your Server Environment** |
|||
<!-- 🔅🔅🔅🔅🔅🔅🔅 Choose one of the following or write your own 🔅🔅🔅🔅🔅🔅🔅--> |
|||
* demo.thingsboard.io |
|||
* cloud.thingsboard.io |
|||
* own setup |
|||
* cloud or local infrastructure or docker deployment |
|||
* ThingsBoard Version |
|||
* OS Name and Version |
|||
|
|||
**Your Client Environment** |
|||
<!-- 🔅🔅🔅🔅🔅🔅🔅 Choose one of the following or write your own 🔅🔅🔅🔅🔅🔅🔅--> |
|||
**Desktop (please complete the following information):** |
|||
|
|||
* OS: [e.g. iOS] |
|||
* Browser [e.g. chrome, safari] |
|||
* Version [e.g. 22] |
|||
|
|||
**Smartphone (please complete the following information):** |
|||
* Device: [e.g. iPhone6] |
|||
* OS: [e.g. iOS8.1] |
|||
* Browser [e.g. stock browser, safari] |
|||
* Version [e.g. 22] |
|||
|
|||
**Your Device** |
|||
|
|||
* Connectivity |
|||
* MQTT |
|||
* HTTP |
|||
* CoAP |
|||
* Gateway |
|||
* Integration: (Specify name) |
|||
* Device vendor and model |
|||
|
|||
**To Reproduce** |
|||
Steps to reproduce the behavior: |
|||
1. Go to '...' |
|||
2. Click on '....' |
|||
3. Scroll down to '....' |
|||
4. See error |
|||
|
|||
**Expected behavior** |
|||
A clear and concise description of what you expected to happen. |
|||
|
|||
**Screenshots** |
|||
If applicable, add screenshots to help explain your problem. |
|||
|
|||
**Additional context** |
|||
Add any other context about the problem here. |
|||
@ -0,0 +1,20 @@ |
|||
--- |
|||
name: Feature request |
|||
about: Suggest an idea for this project |
|||
title: "[Feature Request]" |
|||
labels: feature |
|||
assignees: ikulikov |
|||
|
|||
--- |
|||
|
|||
**Is your feature request related to a problem? Please describe.** |
|||
A clear and concise description of what the problem is. Ex. I'm always frustrated when [...] |
|||
|
|||
**Describe the solution you'd like** |
|||
A clear and concise description of what you want to happen. |
|||
|
|||
**Describe alternatives you've considered** |
|||
A clear and concise description of any alternative solutions or features you've considered. |
|||
|
|||
**Additional context** |
|||
Add any other context or screenshots about the feature request here. |
|||
@ -0,0 +1,25 @@ |
|||
--- |
|||
name: Question |
|||
about: Describe your questions in details |
|||
title: "[Question] your title here" |
|||
labels: question |
|||
assignees: ashvayka |
|||
|
|||
--- |
|||
|
|||
**Component** |
|||
|
|||
<!-- Choose one of the following and delete all others. --> |
|||
* UI |
|||
* Rule Engine |
|||
* Installation |
|||
* Generic |
|||
|
|||
**Description** |
|||
A clear and concise details. |
|||
|
|||
**Environment** |
|||
<!-- Add information about your environment and ThingsBoard version if applicable --> |
|||
* OS: name and version |
|||
* ThingsBoard: version |
|||
* Browser: name and version |
|||
@ -1,174 +0,0 @@ |
|||
/** |
|||
* 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. |
|||
*/ |
|||
import org.apache.tools.ant.filters.ReplaceTokens |
|||
|
|||
buildscript { |
|||
ext { |
|||
osPackageVersion = "3.8.0" |
|||
} |
|||
repositories { |
|||
jcenter() |
|||
} |
|||
dependencies { |
|||
classpath("com.netflix.nebula:gradle-ospackage-plugin:${osPackageVersion}") |
|||
} |
|||
} |
|||
|
|||
apply plugin: "nebula.ospackage" |
|||
|
|||
buildDir = projectBuildDir |
|||
version = projectVersion |
|||
distsDirName = "./" |
|||
|
|||
// OS Package plugin configuration |
|||
ospackage { |
|||
packageName = pkgName |
|||
version = "${project.version}" |
|||
release = 1 |
|||
os = LINUX |
|||
type = BINARY |
|||
|
|||
into pkgInstallFolder |
|||
|
|||
user pkgName |
|||
permissionGroup pkgName |
|||
|
|||
// Copy the actual .jar file |
|||
from(mainJar) { |
|||
// Strip the version from the jar filename |
|||
rename { String fileName -> |
|||
"${pkgName}.jar" |
|||
} |
|||
fileMode 0500 |
|||
into "bin" |
|||
} |
|||
|
|||
// Copy the install files |
|||
from("target/bin/install/install.sh") { |
|||
fileMode 0775 |
|||
into "bin/install" |
|||
} |
|||
|
|||
from("target/bin/install/upgrade.sh") { |
|||
fileMode 0775 |
|||
into "bin/install" |
|||
} |
|||
|
|||
from("target/bin/install/logback.xml") { |
|||
into "bin/install" |
|||
} |
|||
|
|||
// Copy the config files |
|||
from("target/conf") { |
|||
exclude "${pkgName}.conf" |
|||
fileType CONFIG | NOREPLACE |
|||
fileMode 0754 |
|||
into "conf" |
|||
} |
|||
|
|||
// Copy the data files |
|||
from("target/data") { |
|||
fileType CONFIG | NOREPLACE |
|||
fileMode 0754 |
|||
into "data" |
|||
} |
|||
|
|||
// Copy the extensions files |
|||
from("target/extensions") { |
|||
into "extensions" |
|||
} |
|||
} |
|||
|
|||
// Configure our RPM build task |
|||
buildRpm { |
|||
|
|||
arch = NOARCH |
|||
|
|||
version = projectVersion.replace('-', '') |
|||
archiveName = "${pkgName}.rpm" |
|||
|
|||
requires("java-1.8.0") |
|||
|
|||
from("target/conf") { |
|||
include "${pkgName}.conf" |
|||
filter(ReplaceTokens, tokens: ['pkg.platform': 'rpm']) |
|||
fileType CONFIG | NOREPLACE |
|||
fileMode 0754 |
|||
into "${pkgInstallFolder}/conf" |
|||
} |
|||
|
|||
preInstall file("${buildDir}/control/rpm/preinst") |
|||
postInstall file("${buildDir}/control/rpm/postinst") |
|||
preUninstall file("${buildDir}/control/rpm/prerm") |
|||
postUninstall file("${buildDir}/control/rpm/postrm") |
|||
|
|||
user pkgName |
|||
permissionGroup pkgName |
|||
|
|||
// Copy the system unit files |
|||
from("${buildDir}/control/${pkgName}.service") { |
|||
addParentDirs = false |
|||
fileMode 0644 |
|||
into "/usr/lib/systemd/system" |
|||
} |
|||
|
|||
directory(pkgLogFolder, 0755) |
|||
link("${pkgInstallFolder}/bin/${pkgName}.yml", "${pkgInstallFolder}/conf/${pkgName}.yml") |
|||
link("/etc/${pkgName}/conf", "${pkgInstallFolder}/conf") |
|||
} |
|||
|
|||
// Same as the buildRpm task |
|||
buildDeb { |
|||
|
|||
arch = "all" |
|||
|
|||
archiveName = "${pkgName}.deb" |
|||
|
|||
requires("openjdk-8-jre").or("java8-runtime").or("oracle-java8-installer").or("openjdk-8-jre-headless") |
|||
|
|||
from("target/conf") { |
|||
include "${pkgName}.conf" |
|||
filter(ReplaceTokens, tokens: ['pkg.platform': 'deb']) |
|||
fileType CONFIG | NOREPLACE |
|||
fileMode 0754 |
|||
into "${pkgInstallFolder}/conf" |
|||
} |
|||
|
|||
configurationFile("${pkgInstallFolder}/conf/${pkgName}.conf") |
|||
configurationFile("${pkgInstallFolder}/conf/${pkgName}.yml") |
|||
configurationFile("${pkgInstallFolder}/conf/logback.xml") |
|||
configurationFile("${pkgInstallFolder}/conf/actor-system.conf") |
|||
|
|||
preInstall file("${buildDir}/control/deb/preinst") |
|||
postInstall file("${buildDir}/control/deb/postinst") |
|||
preUninstall file("${buildDir}/control/deb/prerm") |
|||
postUninstall file("${buildDir}/control/deb/postrm") |
|||
|
|||
user pkgName |
|||
permissionGroup pkgName |
|||
|
|||
// Copy the system unit files |
|||
from("${buildDir}/control/${pkgName}.service") { |
|||
addParentDirs = false |
|||
fileMode 0644 |
|||
into "/lib/systemd/system" |
|||
} |
|||
|
|||
directory(pkgLogFolder, 0755) |
|||
link("/etc/init.d/${pkgName}", "${pkgInstallFolder}/bin/${pkgName}.jar") |
|||
link("${pkgInstallFolder}/bin/${pkgName}.yml", "${pkgInstallFolder}/conf/${pkgName}.yml") |
|||
link("/etc/${pkgName}/conf", "${pkgInstallFolder}/conf") |
|||
} |
|||
File diff suppressed because it is too large
@ -0,0 +1,520 @@ |
|||
{ |
|||
"title": "Rule Engine Statistics", |
|||
"configuration": { |
|||
"widgets": { |
|||
"81987f19-3eac-e4ce-b790-d96e9b54d9a0": { |
|||
"isSystemType": true, |
|||
"bundleAlias": "charts", |
|||
"typeAlias": "basic_timeseries", |
|||
"type": "timeseries", |
|||
"title": "New widget", |
|||
"sizeX": 12, |
|||
"sizeY": 7, |
|||
"config": { |
|||
"datasources": [ |
|||
{ |
|||
"type": "entity", |
|||
"dataKeys": [ |
|||
{ |
|||
"name": "successfulMsgs", |
|||
"type": "timeseries", |
|||
"label": "${entityName} Successful", |
|||
"color": "#4caf50", |
|||
"settings": { |
|||
"excludeFromStacking": false, |
|||
"hideDataByDefault": false, |
|||
"disableDataHiding": false, |
|||
"removeFromLegend": false, |
|||
"showLines": true, |
|||
"fillLines": false, |
|||
"showPoints": false, |
|||
"showPointShape": "circle", |
|||
"pointShapeFormatter": "var size = radius * Math.sqrt(Math.PI) / 2;\nctx.moveTo(x - size, y - size);\nctx.lineTo(x + size, y + size);\nctx.moveTo(x - size, y + size);\nctx.lineTo(x + size, y - size);", |
|||
"showPointsLineWidth": 5, |
|||
"showPointsRadius": 3, |
|||
"showSeparateAxis": false, |
|||
"axisPosition": "left", |
|||
"thresholds": [ |
|||
{ |
|||
"thresholdValueSource": "predefinedValue" |
|||
} |
|||
], |
|||
"comparisonSettings": { |
|||
"showValuesForComparison": true |
|||
} |
|||
}, |
|||
"_hash": 0.15490750967648736 |
|||
}, |
|||
{ |
|||
"name": "failedMsgs", |
|||
"type": "timeseries", |
|||
"label": "${entityName} Permanent Failures", |
|||
"color": "#ef5350", |
|||
"settings": { |
|||
"excludeFromStacking": false, |
|||
"hideDataByDefault": false, |
|||
"disableDataHiding": false, |
|||
"removeFromLegend": false, |
|||
"showLines": true, |
|||
"fillLines": false, |
|||
"showPoints": false, |
|||
"showPointShape": "circle", |
|||
"pointShapeFormatter": "var size = radius * Math.sqrt(Math.PI) / 2;\nctx.moveTo(x - size, y - size);\nctx.lineTo(x + size, y + size);\nctx.moveTo(x - size, y + size);\nctx.lineTo(x + size, y - size);", |
|||
"showPointsLineWidth": 5, |
|||
"showPointsRadius": 3, |
|||
"showSeparateAxis": false, |
|||
"axisPosition": "left", |
|||
"thresholds": [ |
|||
{ |
|||
"thresholdValueSource": "predefinedValue" |
|||
} |
|||
], |
|||
"comparisonSettings": { |
|||
"showValuesForComparison": true |
|||
} |
|||
}, |
|||
"_hash": 0.4186621166514697 |
|||
}, |
|||
{ |
|||
"name": "tmpFailed", |
|||
"type": "timeseries", |
|||
"label": "${entityName} Processing Failures", |
|||
"color": "#ffc107", |
|||
"settings": { |
|||
"excludeFromStacking": false, |
|||
"hideDataByDefault": false, |
|||
"disableDataHiding": false, |
|||
"removeFromLegend": false, |
|||
"showLines": true, |
|||
"fillLines": false, |
|||
"showPoints": false, |
|||
"showPointShape": "circle", |
|||
"pointShapeFormatter": "var size = radius * Math.sqrt(Math.PI) / 2;\nctx.moveTo(x - size, y - size);\nctx.lineTo(x + size, y + size);\nctx.moveTo(x - size, y + size);\nctx.lineTo(x + size, y - size);", |
|||
"showPointsLineWidth": 5, |
|||
"showPointsRadius": 3, |
|||
"showSeparateAxis": false, |
|||
"axisPosition": "left", |
|||
"thresholds": [ |
|||
{ |
|||
"thresholdValueSource": "predefinedValue" |
|||
} |
|||
], |
|||
"comparisonSettings": { |
|||
"showValuesForComparison": true |
|||
} |
|||
}, |
|||
"_hash": 0.49891007198715376 |
|||
} |
|||
], |
|||
"entityAliasId": "140f23dd-e3a0-ed98-6189-03c49d2d8018" |
|||
} |
|||
], |
|||
"timewindow": { |
|||
"realtime": { |
|||
"interval": 1000, |
|||
"timewindowMs": 300000 |
|||
}, |
|||
"aggregation": { |
|||
"type": "NONE", |
|||
"limit": 8640 |
|||
}, |
|||
"hideInterval": false, |
|||
"hideAggregation": false, |
|||
"hideAggInterval": false |
|||
}, |
|||
"showTitle": true, |
|||
"backgroundColor": "#fff", |
|||
"color": "rgba(0, 0, 0, 0.87)", |
|||
"padding": "8px", |
|||
"settings": { |
|||
"shadowSize": 4, |
|||
"fontColor": "#545454", |
|||
"fontSize": 10, |
|||
"xaxis": { |
|||
"showLabels": true, |
|||
"color": "#545454" |
|||
}, |
|||
"yaxis": { |
|||
"showLabels": true, |
|||
"color": "#545454" |
|||
}, |
|||
"grid": { |
|||
"color": "#545454", |
|||
"tickColor": "#DDDDDD", |
|||
"verticalLines": true, |
|||
"horizontalLines": true, |
|||
"outlineWidth": 1 |
|||
}, |
|||
"stack": false, |
|||
"tooltipIndividual": false, |
|||
"timeForComparison": "months", |
|||
"xaxisSecond": { |
|||
"axisPosition": "top", |
|||
"showLabels": true |
|||
} |
|||
}, |
|||
"title": "Queue Stats", |
|||
"dropShadow": true, |
|||
"enableFullscreen": true, |
|||
"titleStyle": { |
|||
"fontSize": "16px", |
|||
"fontWeight": 400 |
|||
}, |
|||
"mobileHeight": null, |
|||
"showTitleIcon": false, |
|||
"titleIcon": null, |
|||
"iconColor": "rgba(0, 0, 0, 0.87)", |
|||
"iconSize": "24px", |
|||
"titleTooltip": "", |
|||
"widgetStyle": {}, |
|||
"useDashboardTimewindow": false, |
|||
"displayTimewindow": true, |
|||
"showLegend": true, |
|||
"actions": {}, |
|||
"legendConfig": { |
|||
"direction": "column", |
|||
"position": "bottom", |
|||
"showMin": true, |
|||
"showMax": true, |
|||
"showAvg": false, |
|||
"showTotal": true |
|||
} |
|||
}, |
|||
"id": "81987f19-3eac-e4ce-b790-d96e9b54d9a0" |
|||
}, |
|||
"5eb79712-5c24-3060-7e4f-6af36b8f842d": { |
|||
"isSystemType": true, |
|||
"bundleAlias": "cards", |
|||
"typeAlias": "timeseries_table", |
|||
"type": "timeseries", |
|||
"title": "New widget", |
|||
"sizeX": 24, |
|||
"sizeY": 5, |
|||
"config": { |
|||
"datasources": [ |
|||
{ |
|||
"type": "entity", |
|||
"dataKeys": [ |
|||
{ |
|||
"name": "ruleEngineException", |
|||
"type": "timeseries", |
|||
"label": "Rule Chain", |
|||
"color": "#2196f3", |
|||
"settings": { |
|||
"useCellStyleFunction": false, |
|||
"useCellContentFunction": true, |
|||
"cellContentFunction": "return JSON.parse(value).ruleChainName;" |
|||
}, |
|||
"_hash": 0.9954481282345906 |
|||
}, |
|||
{ |
|||
"name": "ruleEngineException", |
|||
"type": "timeseries", |
|||
"label": "Rule Node", |
|||
"color": "#4caf50", |
|||
"settings": { |
|||
"useCellStyleFunction": false, |
|||
"useCellContentFunction": true, |
|||
"cellContentFunction": "return JSON.parse(value).ruleNodeName;" |
|||
}, |
|||
"_hash": 0.18580357036589978 |
|||
}, |
|||
{ |
|||
"name": "ruleEngineException", |
|||
"type": "timeseries", |
|||
"label": "Latest Error", |
|||
"color": "#f44336", |
|||
"settings": { |
|||
"useCellStyleFunction": false, |
|||
"useCellContentFunction": true, |
|||
"cellContentFunction": "return JSON.parse(value).message;" |
|||
}, |
|||
"_hash": 0.7255162989552142 |
|||
} |
|||
], |
|||
"entityAliasId": "140f23dd-e3a0-ed98-6189-03c49d2d8018" |
|||
} |
|||
], |
|||
"timewindow": { |
|||
"realtime": { |
|||
"interval": 1000, |
|||
"timewindowMs": 86400000 |
|||
}, |
|||
"aggregation": { |
|||
"type": "NONE", |
|||
"limit": 200 |
|||
} |
|||
}, |
|||
"showTitle": true, |
|||
"backgroundColor": "rgb(255, 255, 255)", |
|||
"color": "rgba(0, 0, 0, 0.87)", |
|||
"padding": "8px", |
|||
"settings": { |
|||
"showTimestamp": true, |
|||
"displayPagination": true, |
|||
"defaultPageSize": 10 |
|||
}, |
|||
"title": "Exceptions", |
|||
"dropShadow": true, |
|||
"enableFullscreen": true, |
|||
"titleStyle": { |
|||
"fontSize": "16px", |
|||
"fontWeight": 400 |
|||
}, |
|||
"useDashboardTimewindow": false, |
|||
"showLegend": false, |
|||
"widgetStyle": {}, |
|||
"actions": {}, |
|||
"showTitleIcon": false, |
|||
"titleIcon": null, |
|||
"iconColor": "rgba(0, 0, 0, 0.87)", |
|||
"iconSize": "24px", |
|||
"titleTooltip": "", |
|||
"displayTimewindow": true |
|||
}, |
|||
"id": "5eb79712-5c24-3060-7e4f-6af36b8f842d" |
|||
}, |
|||
"ad3f1417-87a8-750e-fc67-49a2de1466d4": { |
|||
"isSystemType": true, |
|||
"bundleAlias": "charts", |
|||
"typeAlias": "basic_timeseries", |
|||
"type": "timeseries", |
|||
"title": "New widget", |
|||
"sizeX": 12, |
|||
"sizeY": 7, |
|||
"config": { |
|||
"datasources": [ |
|||
{ |
|||
"type": "entity", |
|||
"dataKeys": [ |
|||
{ |
|||
"name": "timeoutMsgs", |
|||
"type": "timeseries", |
|||
"label": "${entityName} Permanent Timeouts", |
|||
"color": "#4caf50", |
|||
"settings": { |
|||
"excludeFromStacking": false, |
|||
"hideDataByDefault": false, |
|||
"disableDataHiding": false, |
|||
"removeFromLegend": false, |
|||
"showLines": true, |
|||
"fillLines": false, |
|||
"showPoints": false, |
|||
"showPointShape": "circle", |
|||
"pointShapeFormatter": "var size = radius * Math.sqrt(Math.PI) / 2;\nctx.moveTo(x - size, y - size);\nctx.lineTo(x + size, y + size);\nctx.moveTo(x - size, y + size);\nctx.lineTo(x + size, y - size);", |
|||
"showPointsLineWidth": 5, |
|||
"showPointsRadius": 3, |
|||
"showSeparateAxis": false, |
|||
"axisPosition": "left", |
|||
"thresholds": [ |
|||
{ |
|||
"thresholdValueSource": "predefinedValue" |
|||
} |
|||
], |
|||
"comparisonSettings": { |
|||
"showValuesForComparison": true |
|||
} |
|||
}, |
|||
"_hash": 0.565222981550328 |
|||
}, |
|||
{ |
|||
"name": "tmpTimeout", |
|||
"type": "timeseries", |
|||
"label": "${entityName} Processing Timeouts", |
|||
"color": "#9c27b0", |
|||
"settings": { |
|||
"excludeFromStacking": false, |
|||
"hideDataByDefault": false, |
|||
"disableDataHiding": false, |
|||
"removeFromLegend": false, |
|||
"showLines": true, |
|||
"fillLines": false, |
|||
"showPoints": false, |
|||
"showPointShape": "circle", |
|||
"pointShapeFormatter": "var size = radius * Math.sqrt(Math.PI) / 2;\nctx.moveTo(x - size, y - size);\nctx.lineTo(x + size, y + size);\nctx.moveTo(x - size, y + size);\nctx.lineTo(x + size, y - size);", |
|||
"showPointsLineWidth": 5, |
|||
"showPointsRadius": 3, |
|||
"showSeparateAxis": false, |
|||
"axisPosition": "left", |
|||
"thresholds": [ |
|||
{ |
|||
"thresholdValueSource": "predefinedValue" |
|||
} |
|||
], |
|||
"comparisonSettings": { |
|||
"showValuesForComparison": true |
|||
} |
|||
}, |
|||
"_hash": 0.2679547062508352 |
|||
} |
|||
], |
|||
"entityAliasId": "140f23dd-e3a0-ed98-6189-03c49d2d8018" |
|||
} |
|||
], |
|||
"timewindow": { |
|||
"realtime": { |
|||
"interval": 1000, |
|||
"timewindowMs": 300000 |
|||
}, |
|||
"aggregation": { |
|||
"type": "NONE", |
|||
"limit": 8640 |
|||
}, |
|||
"hideInterval": false, |
|||
"hideAggregation": false, |
|||
"hideAggInterval": false |
|||
}, |
|||
"showTitle": true, |
|||
"backgroundColor": "#fff", |
|||
"color": "rgba(0, 0, 0, 0.87)", |
|||
"padding": "8px", |
|||
"settings": { |
|||
"shadowSize": 4, |
|||
"fontColor": "#545454", |
|||
"fontSize": 10, |
|||
"xaxis": { |
|||
"showLabels": true, |
|||
"color": "#545454" |
|||
}, |
|||
"yaxis": { |
|||
"showLabels": true, |
|||
"color": "#545454" |
|||
}, |
|||
"grid": { |
|||
"color": "#545454", |
|||
"tickColor": "#DDDDDD", |
|||
"verticalLines": true, |
|||
"horizontalLines": true, |
|||
"outlineWidth": 1 |
|||
}, |
|||
"stack": false, |
|||
"tooltipIndividual": false, |
|||
"timeForComparison": "months", |
|||
"xaxisSecond": { |
|||
"axisPosition": "top", |
|||
"showLabels": true |
|||
} |
|||
}, |
|||
"title": "Processing Failures and Timeouts", |
|||
"dropShadow": true, |
|||
"enableFullscreen": true, |
|||
"titleStyle": { |
|||
"fontSize": "16px", |
|||
"fontWeight": 400 |
|||
}, |
|||
"mobileHeight": null, |
|||
"showTitleIcon": false, |
|||
"titleIcon": null, |
|||
"iconColor": "rgba(0, 0, 0, 0.87)", |
|||
"iconSize": "24px", |
|||
"titleTooltip": "", |
|||
"widgetStyle": {}, |
|||
"useDashboardTimewindow": false, |
|||
"displayTimewindow": true, |
|||
"showLegend": true, |
|||
"actions": {}, |
|||
"legendConfig": { |
|||
"direction": "column", |
|||
"position": "bottom", |
|||
"showMin": true, |
|||
"showMax": true, |
|||
"showAvg": false, |
|||
"showTotal": true |
|||
} |
|||
}, |
|||
"id": "ad3f1417-87a8-750e-fc67-49a2de1466d4" |
|||
} |
|||
}, |
|||
"states": { |
|||
"default": { |
|||
"name": "Rule Engine Statistics", |
|||
"root": true, |
|||
"layouts": { |
|||
"main": { |
|||
"widgets": { |
|||
"81987f19-3eac-e4ce-b790-d96e9b54d9a0": { |
|||
"sizeX": 12, |
|||
"sizeY": 7, |
|||
"mobileHeight": null, |
|||
"row": 0, |
|||
"col": 0 |
|||
}, |
|||
"5eb79712-5c24-3060-7e4f-6af36b8f842d": { |
|||
"sizeX": 24, |
|||
"sizeY": 5, |
|||
"row": 7, |
|||
"col": 0 |
|||
}, |
|||
"ad3f1417-87a8-750e-fc67-49a2de1466d4": { |
|||
"sizeX": 12, |
|||
"sizeY": 7, |
|||
"mobileHeight": null, |
|||
"row": 0, |
|||
"col": 12 |
|||
} |
|||
}, |
|||
"gridSettings": { |
|||
"backgroundColor": "#eeeeee", |
|||
"color": "rgba(0,0,0,0.870588)", |
|||
"columns": 24, |
|||
"margins": [ |
|||
10, |
|||
10 |
|||
], |
|||
"backgroundSizeMode": "100%", |
|||
"autoFillHeight": true, |
|||
"mobileAutoFillHeight": false, |
|||
"mobileRowHeight": 70 |
|||
} |
|||
} |
|||
} |
|||
} |
|||
}, |
|||
"entityAliases": { |
|||
"140f23dd-e3a0-ed98-6189-03c49d2d8018": { |
|||
"id": "140f23dd-e3a0-ed98-6189-03c49d2d8018", |
|||
"alias": "TbServiceQueues", |
|||
"filter": { |
|||
"type": "assetType", |
|||
"resolveMultiple": true, |
|||
"assetType": "TbServiceQueue", |
|||
"assetNameFilter": "" |
|||
} |
|||
} |
|||
}, |
|||
"timewindow": { |
|||
"displayValue": "", |
|||
"selectedTab": 0, |
|||
"hideInterval": false, |
|||
"hideAggregation": false, |
|||
"hideAggInterval": false, |
|||
"realtime": { |
|||
"interval": 1000, |
|||
"timewindowMs": 60000 |
|||
}, |
|||
"history": { |
|||
"historyType": 0, |
|||
"interval": 1000, |
|||
"timewindowMs": 60000, |
|||
"fixedTimewindow": { |
|||
"startTimeMs": 1586176634823, |
|||
"endTimeMs": 1586263034823 |
|||
} |
|||
}, |
|||
"aggregation": { |
|||
"type": "AVG", |
|||
"limit": 25000 |
|||
} |
|||
}, |
|||
"settings": { |
|||
"stateControllerId": "entity", |
|||
"showTitle": false, |
|||
"showDashboardsSelect": true, |
|||
"showEntitiesSelect": true, |
|||
"showDashboardTimewindow": true, |
|||
"showDashboardExport": true, |
|||
"toolbarAlwaysOpen": true |
|||
} |
|||
}, |
|||
"name": "Rule Engine Statistics" |
|||
} |
|||
File diff suppressed because it is too large
@ -0,0 +1,188 @@ |
|||
{ |
|||
"ruleChain": { |
|||
"additionalInfo": null, |
|||
"name": "Root Rule Chain", |
|||
"firstRuleNodeId": null, |
|||
"root": true, |
|||
"debugMode": false, |
|||
"configuration": null |
|||
}, |
|||
"metadata": { |
|||
"firstNodeIndex": 3, |
|||
"nodes": [ |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 1069, |
|||
"layoutY": 267 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.filter.TbJsFilterNode", |
|||
"name": "Is Thermostat?", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return msg.id.entityType === \"DEVICE\" && msg.type === \"thermostat\";" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 824, |
|||
"layoutY": 156 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode", |
|||
"name": "Save Timeseries", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"defaultTTL": 0 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 825, |
|||
"layoutY": 52 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode", |
|||
"name": "Save Client Attributes", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"scope": "CLIENT_SCOPE" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 347, |
|||
"layoutY": 149 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode", |
|||
"name": "Message Type Switch", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"version": 0 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 839, |
|||
"layoutY": 345 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbLogNode", |
|||
"name": "Log RPC from Device", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 832, |
|||
"layoutY": 407 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbLogNode", |
|||
"name": "Log Other", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 825, |
|||
"layoutY": 468 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode", |
|||
"name": "RPC Call Request", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"timeoutInSeconds": 60 |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 1069, |
|||
"layoutY": 90 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.filter.TbJsFilterNode", |
|||
"name": "Is Thermostat?", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "return metadata[\"deviceType\"] === \"thermostat\";" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 1090, |
|||
"layoutY": 360 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbCreateRelationNode", |
|||
"name": "Relate to Asset", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"direction": "FROM", |
|||
"relationType": "ToAlarmPropagationAsset", |
|||
"entityType": "ASSET", |
|||
"entityNamePattern": "Thermostat Alarms", |
|||
"entityTypePattern": "AlarmPropagationAsset", |
|||
"entityCacheExpiration": 300, |
|||
"createEntityIfNotExists": true, |
|||
"changeOriginatorToRelatedEntity": false, |
|||
"removeCurrentRelations": false |
|||
} |
|||
} |
|||
], |
|||
"connections": [ |
|||
{ |
|||
"fromIndex": 0, |
|||
"toIndex": 8, |
|||
"type": "True" |
|||
}, |
|||
{ |
|||
"fromIndex": 1, |
|||
"toIndex": 7, |
|||
"type": "Success" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 5, |
|||
"type": "Other" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 2, |
|||
"type": "Post attributes" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 1, |
|||
"type": "Post telemetry" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 4, |
|||
"type": "RPC Request from Device" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 6, |
|||
"type": "RPC Request to Device" |
|||
}, |
|||
{ |
|||
"fromIndex": 3, |
|||
"toIndex": 0, |
|||
"type": "Entity Created" |
|||
} |
|||
], |
|||
"ruleChainConnections": [ |
|||
{ |
|||
"fromIndex": 7, |
|||
"targetRuleChainId": { |
|||
"entityType": "RULE_CHAIN", |
|||
"id": "25e26570-89ed-11ea-a650-cd6e14e633bd" |
|||
}, |
|||
"additionalInfo": { |
|||
"layoutX": 1109, |
|||
"layoutY": 182, |
|||
"ruleChainNodeId": "rule-chain-node-10" |
|||
}, |
|||
"type": "True" |
|||
} |
|||
] |
|||
} |
|||
} |
|||
@ -0,0 +1,141 @@ |
|||
{ |
|||
"ruleChain": { |
|||
"additionalInfo": null, |
|||
"name": "Thermostat Alarms", |
|||
"firstRuleNodeId": null, |
|||
"root": false, |
|||
"debugMode": false, |
|||
"configuration": null |
|||
}, |
|||
"metadata": { |
|||
"firstNodeIndex": 5, |
|||
"nodes": [ |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 929, |
|||
"layoutY": 67 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbCreateAlarmNode", |
|||
"name": "Create Temp Alarm", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"alarmType": "High Temperature", |
|||
"alarmDetailsBuildJs": "var details = {};\nif (metadata.prevAlarmDetails) {\n details = JSON.parse(metadata.prevAlarmDetails);\n}\ndetails.triggerValue = msg.temperature;\nreturn details;", |
|||
"severity": "MAJOR", |
|||
"propagate": true, |
|||
"useMessageAlarmData": false, |
|||
"relationTypes": [ |
|||
"ToAlarmPropagationAsset" |
|||
] |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 930, |
|||
"layoutY": 201 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbClearAlarmNode", |
|||
"name": "Clear Temp Alarm", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"alarmType": "High Temperature", |
|||
"alarmDetailsBuildJs": "var details = {};\nif (metadata.prevAlarmDetails) {\n details = JSON.parse(metadata.prevAlarmDetails);\n}\nreturn details;" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 930, |
|||
"layoutY": 131 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbCreateAlarmNode", |
|||
"name": "Create Humidity Alarm", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"alarmType": "Low Humidity", |
|||
"alarmDetailsBuildJs": "var details = {};\nif (metadata.prevAlarmDetails) {\n details = JSON.parse(metadata.prevAlarmDetails);\n}\ndetails.triggerValue = msg.humidity;\nreturn details;", |
|||
"severity": "MINOR", |
|||
"propagate": true, |
|||
"useMessageAlarmData": false, |
|||
"relationTypes": [ |
|||
"ToAlarmPropagationAsset" |
|||
] |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 929, |
|||
"layoutY": 275 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.action.TbClearAlarmNode", |
|||
"name": "Clear Humidity Alarm", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"alarmType": "Low Humidity", |
|||
"alarmDetailsBuildJs": "var details = {};\nif (metadata.prevAlarmDetails) {\n details = JSON.parse(metadata.prevAlarmDetails);\n}\nreturn details;" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 586, |
|||
"layoutY": 148 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.filter.TbJsSwitchNode", |
|||
"name": "Check Alarms", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"jsScript": "var relations = [];\nif(metadata[\"ss_alarmTemperature\"] === \"true\"){\n if(msg.temperature > metadata[\"ss_thresholdTemperature\"]){\n relations.push(\"NewTempAlarm\");\n } else {\n relations.push(\"ClearTempAlarm\");\n }\n}\nif(metadata[\"ss_alarmHumidity\"] === \"true\"){\n if(msg.humidity < metadata[\"ss_thresholdHumidity\"]){\n relations.push(\"NewHumidityAlarm\");\n } else {\n relations.push(\"ClearHumidityAlarm\");\n }\n}\n\nreturn relations;" |
|||
} |
|||
}, |
|||
{ |
|||
"additionalInfo": { |
|||
"layoutX": 321, |
|||
"layoutY": 149 |
|||
}, |
|||
"type": "org.thingsboard.rule.engine.metadata.TbGetAttributesNode", |
|||
"name": "Fetch Configuration", |
|||
"debugMode": false, |
|||
"configuration": { |
|||
"clientAttributeNames": [], |
|||
"sharedAttributeNames": [], |
|||
"serverAttributeNames": [ |
|||
"alarmTemperature", |
|||
"thresholdTemperature", |
|||
"alarmHumidity", |
|||
"thresholdHumidity" |
|||
], |
|||
"latestTsKeyNames": [], |
|||
"tellFailureIfAbsent": false, |
|||
"getLatestValueWithTs": false |
|||
} |
|||
} |
|||
], |
|||
"connections": [ |
|||
{ |
|||
"fromIndex": 4, |
|||
"toIndex": 0, |
|||
"type": "NewTempAlarm" |
|||
}, |
|||
{ |
|||
"fromIndex": 4, |
|||
"toIndex": 1, |
|||
"type": "ClearTempAlarm" |
|||
}, |
|||
{ |
|||
"fromIndex": 4, |
|||
"toIndex": 2, |
|||
"type": "NewHumidityAlarm" |
|||
}, |
|||
{ |
|||
"fromIndex": 4, |
|||
"toIndex": 3, |
|||
"type": "ClearHumidityAlarm" |
|||
}, |
|||
{ |
|||
"fromIndex": 5, |
|||
"toIndex": 4, |
|||
"type": "Success" |
|||
} |
|||
], |
|||
"ruleChainConnections": null |
|||
} |
|||
} |
|||
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@ -0,0 +1,123 @@ |
|||
-- |
|||
-- 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. |
|||
-- |
|||
|
|||
CREATE OR REPLACE PROCEDURE drop_partitions_by_max_ttl(IN partition_type varchar, IN system_ttl bigint, INOUT deleted bigint) |
|||
LANGUAGE plpgsql AS |
|||
$$ |
|||
DECLARE |
|||
max_tenant_ttl bigint; |
|||
max_customer_ttl bigint; |
|||
max_ttl bigint; |
|||
date timestamp; |
|||
partition_by_max_ttl_date varchar; |
|||
partition_month varchar; |
|||
partition_day varchar; |
|||
partition_year varchar; |
|||
partition varchar; |
|||
partition_to_delete varchar; |
|||
|
|||
|
|||
BEGIN |
|||
SELECT max(attribute_kv.long_v) |
|||
FROM tenant |
|||
INNER JOIN attribute_kv ON tenant.id = attribute_kv.entity_id |
|||
WHERE attribute_kv.attribute_key = 'TTL' |
|||
into max_tenant_ttl; |
|||
SELECT max(attribute_kv.long_v) |
|||
FROM customer |
|||
INNER JOIN attribute_kv ON customer.id = attribute_kv.entity_id |
|||
WHERE attribute_kv.attribute_key = 'TTL' |
|||
into max_customer_ttl; |
|||
max_ttl := GREATEST(system_ttl, max_customer_ttl, max_tenant_ttl); |
|||
if max_ttl IS NOT NULL AND max_ttl > 0 THEN |
|||
date := to_timestamp(EXTRACT(EPOCH FROM current_timestamp) - (max_ttl / 1000)); |
|||
partition_by_max_ttl_date := get_partition_by_max_ttl_date(partition_type, date); |
|||
RAISE NOTICE 'Partition by max ttl: %', partition_by_max_ttl_date; |
|||
IF partition_by_max_ttl_date IS NOT NULL THEN |
|||
CASE |
|||
WHEN partition_type = 'DAYS' THEN |
|||
partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); |
|||
partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); |
|||
partition_day := SPLIT_PART(partition_by_max_ttl_date, '_', 5); |
|||
WHEN partition_type = 'MONTHS' THEN |
|||
partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); |
|||
partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); |
|||
ELSE |
|||
partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); |
|||
END CASE; |
|||
FOR partition IN SELECT tablename |
|||
FROM pg_tables |
|||
WHERE schemaname = 'public' |
|||
AND tablename like 'ts_kv_' || '%' |
|||
AND tablename != 'ts_kv_latest' |
|||
AND tablename != 'ts_kv_dictionary' |
|||
LOOP |
|||
IF partition != partition_by_max_ttl_date THEN |
|||
IF partition_year IS NOT NULL THEN |
|||
IF SPLIT_PART(partition, '_', 3)::integer < partition_year::integer THEN |
|||
partition_to_delete := partition; |
|||
ELSE |
|||
IF partition_month IS NOT NULL THEN |
|||
IF SPLIT_PART(partition, '_', 4)::integer < partition_month::integer THEN |
|||
partition_to_delete := partition; |
|||
ELSE |
|||
IF partition_day IS NOT NULL THEN |
|||
IF SPLIT_PART(partition, '_', 5)::integer < partition_day::integer THEN |
|||
partition_to_delete := partition; |
|||
END IF; |
|||
END IF; |
|||
END IF; |
|||
END IF; |
|||
END IF; |
|||
END IF; |
|||
END IF; |
|||
IF partition_to_delete IS NOT NULL THEN |
|||
RAISE NOTICE 'Partition to delete by max ttl: %', partition_to_delete; |
|||
EXECUTE format('DROP TABLE %I', partition_to_delete); |
|||
deleted := deleted + 1; |
|||
END IF; |
|||
END LOOP; |
|||
END IF; |
|||
END IF; |
|||
END |
|||
$$; |
|||
|
|||
CREATE OR REPLACE FUNCTION get_partition_by_max_ttl_date(IN partition_type varchar, IN date timestamp, OUT partition varchar) AS |
|||
$$ |
|||
BEGIN |
|||
CASE |
|||
WHEN partition_type = 'DAYS' THEN |
|||
partition := 'ts_kv_' || to_char(date, 'yyyy') || '_' || to_char(date, 'MM') || '_' || to_char(date, 'dd'); |
|||
WHEN partition_type = 'MONTHS' THEN |
|||
partition := 'ts_kv_' || to_char(date, 'yyyy') || '_' || to_char(date, 'MM'); |
|||
WHEN partition_type = 'YEARS' THEN |
|||
partition := 'ts_kv_' || to_char(date, 'yyyy'); |
|||
WHEN partition_type = 'INDEFINITE' THEN |
|||
partition := NULL; |
|||
ELSE |
|||
partition := NULL; |
|||
END CASE; |
|||
IF partition IS NOT NULL THEN |
|||
IF NOT EXISTS(SELECT |
|||
FROM pg_tables |
|||
WHERE schemaname = 'public' |
|||
AND tablename = partition) THEN |
|||
partition := NULL; |
|||
RAISE NOTICE 'Failed to found partition by ttl'; |
|||
END IF; |
|||
END IF; |
|||
END; |
|||
$$ LANGUAGE plpgsql; |
|||
@ -0,0 +1,150 @@ |
|||
-- |
|||
-- 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. |
|||
-- |
|||
|
|||
CREATE OR REPLACE FUNCTION to_uuid(IN entity_id varchar, OUT uuid_id uuid) AS |
|||
$$ |
|||
BEGIN |
|||
uuid_id := substring(entity_id, 8, 8) || '-' || substring(entity_id, 4, 4) || '-1' || substring(entity_id, 1, 3) || |
|||
'-' || substring(entity_id, 16, 4) || '-' || substring(entity_id, 20, 12); |
|||
END; |
|||
$$ LANGUAGE plpgsql; |
|||
|
|||
CREATE OR REPLACE FUNCTION delete_device_records_from_ts_kv(tenant_id varchar, customer_id varchar, ttl bigint, |
|||
OUT deleted bigint) AS |
|||
$$ |
|||
BEGIN |
|||
EXECUTE format( |
|||
'WITH deleted AS (DELETE FROM ts_kv WHERE entity_id IN (SELECT to_uuid(device.id) as entity_id FROM device WHERE tenant_id = %L and customer_id = %L) AND ts < %L::bigint RETURNING *) SELECT count(*) FROM deleted', |
|||
tenant_id, customer_id, ttl) into deleted; |
|||
END; |
|||
$$ LANGUAGE plpgsql; |
|||
|
|||
CREATE OR REPLACE FUNCTION delete_asset_records_from_ts_kv(tenant_id varchar, customer_id varchar, ttl bigint, |
|||
OUT deleted bigint) AS |
|||
$$ |
|||
BEGIN |
|||
EXECUTE format( |
|||
'WITH deleted AS (DELETE FROM ts_kv WHERE entity_id IN (SELECT to_uuid(asset.id) as entity_id FROM asset WHERE tenant_id = %L and customer_id = %L) AND ts < %L::bigint RETURNING *) SELECT count(*) FROM deleted', |
|||
tenant_id, customer_id, ttl) into deleted; |
|||
END; |
|||
$$ LANGUAGE plpgsql; |
|||
|
|||
CREATE OR REPLACE FUNCTION delete_customer_records_from_ts_kv(tenant_id varchar, customer_id varchar, ttl bigint, |
|||
OUT deleted bigint) AS |
|||
$$ |
|||
BEGIN |
|||
EXECUTE format( |
|||
'WITH deleted AS (DELETE FROM ts_kv WHERE entity_id IN (SELECT to_uuid(customer.id) as entity_id FROM customer WHERE tenant_id = %L and id = %L) AND ts < %L::bigint RETURNING *) SELECT count(*) FROM deleted', |
|||
tenant_id, customer_id, ttl) into deleted; |
|||
END; |
|||
$$ LANGUAGE plpgsql; |
|||
|
|||
CREATE OR REPLACE PROCEDURE cleanup_timeseries_by_ttl(IN null_uuid varchar(31), |
|||
IN system_ttl bigint, INOUT deleted bigint) |
|||
LANGUAGE plpgsql AS |
|||
$$ |
|||
DECLARE |
|||
tenant_cursor CURSOR FOR select tenant.id as tenant_id |
|||
from tenant; |
|||
tenant_id_record varchar; |
|||
customer_id_record varchar; |
|||
tenant_ttl bigint; |
|||
customer_ttl bigint; |
|||
deleted_for_entities bigint; |
|||
tenant_ttl_ts bigint; |
|||
customer_ttl_ts bigint; |
|||
BEGIN |
|||
OPEN tenant_cursor; |
|||
FETCH tenant_cursor INTO tenant_id_record; |
|||
WHILE FOUND |
|||
LOOP |
|||
EXECUTE format( |
|||
'select attribute_kv.long_v from attribute_kv where attribute_kv.entity_id = %L and attribute_kv.attribute_key = %L', |
|||
tenant_id_record, 'TTL') INTO tenant_ttl; |
|||
if tenant_ttl IS NULL THEN |
|||
tenant_ttl := system_ttl; |
|||
END IF; |
|||
IF tenant_ttl > 0 THEN |
|||
tenant_ttl_ts := (EXTRACT(EPOCH FROM current_timestamp) * 1000 - tenant_ttl::bigint * 1000)::bigint; |
|||
deleted_for_entities := delete_device_records_from_ts_kv(tenant_id_record, null_uuid, tenant_ttl_ts); |
|||
deleted := deleted + deleted_for_entities; |
|||
RAISE NOTICE '% telemetry removed for devices where tenant_id = %', deleted_for_entities, tenant_id_record; |
|||
deleted_for_entities := delete_asset_records_from_ts_kv(tenant_id_record, null_uuid, tenant_ttl_ts); |
|||
deleted := deleted + deleted_for_entities; |
|||
RAISE NOTICE '% telemetry removed for assets where tenant_id = %', deleted_for_entities, tenant_id_record; |
|||
END IF; |
|||
FOR customer_id_record IN |
|||
SELECT customer.id AS customer_id FROM customer WHERE customer.tenant_id = tenant_id_record |
|||
LOOP |
|||
EXECUTE format( |
|||
'select attribute_kv.long_v from attribute_kv where attribute_kv.entity_id = %L and attribute_kv.attribute_key = %L', |
|||
customer_id_record, 'TTL') INTO customer_ttl; |
|||
IF customer_ttl IS NULL THEN |
|||
customer_ttl_ts := tenant_ttl_ts; |
|||
ELSE |
|||
IF customer_ttl > 0 THEN |
|||
customer_ttl_ts := |
|||
(EXTRACT(EPOCH FROM current_timestamp) * 1000 - |
|||
customer_ttl::bigint * 1000)::bigint; |
|||
END IF; |
|||
END IF; |
|||
IF customer_ttl_ts IS NOT NULL AND customer_ttl_ts > 0 THEN |
|||
deleted_for_entities := |
|||
delete_customer_records_from_ts_kv(tenant_id_record, customer_id_record, |
|||
customer_ttl_ts); |
|||
deleted := deleted + deleted_for_entities; |
|||
RAISE NOTICE '% telemetry removed for customer with id = % where tenant_id = %', deleted_for_entities, customer_id_record, tenant_id_record; |
|||
deleted_for_entities := |
|||
delete_device_records_from_ts_kv(tenant_id_record, customer_id_record, |
|||
customer_ttl_ts); |
|||
deleted := deleted + deleted_for_entities; |
|||
RAISE NOTICE '% telemetry removed for devices where tenant_id = % and customer_id = %', deleted_for_entities, tenant_id_record, customer_id_record; |
|||
deleted_for_entities := delete_asset_records_from_ts_kv(tenant_id_record, |
|||
customer_id_record, |
|||
customer_ttl_ts); |
|||
deleted := deleted + deleted_for_entities; |
|||
RAISE NOTICE '% telemetry removed for assets where tenant_id = % and customer_id = %', deleted_for_entities, tenant_id_record, customer_id_record; |
|||
END IF; |
|||
END LOOP; |
|||
FETCH tenant_cursor INTO tenant_id_record; |
|||
END LOOP; |
|||
END |
|||
$$; |
|||
|
|||
CREATE OR REPLACE PROCEDURE cleanup_events_by_ttl(IN ttl bigint, IN debug_ttl bigint, INOUT deleted bigint) |
|||
LANGUAGE plpgsql AS |
|||
$$ |
|||
DECLARE |
|||
ttl_ts bigint; |
|||
debug_ttl_ts bigint; |
|||
ttl_deleted_count bigint DEFAULT 0; |
|||
debug_ttl_deleted_count bigint DEFAULT 0; |
|||
BEGIN |
|||
IF ttl > 0 THEN |
|||
ttl_ts := (EXTRACT(EPOCH FROM current_timestamp) * 1000 - ttl::bigint * 1000)::bigint; |
|||
EXECUTE format( |
|||
'WITH deleted AS (DELETE FROM event WHERE ts < %L::bigint AND (event_type != %L::varchar AND event_type != %L::varchar) RETURNING *) SELECT count(*) FROM deleted', ttl_ts, 'DEBUG_RULE_NODE', 'DEBUG_RULE_CHAIN') into ttl_deleted_count; |
|||
END IF; |
|||
IF debug_ttl > 0 THEN |
|||
debug_ttl_ts := (EXTRACT(EPOCH FROM current_timestamp) * 1000 - debug_ttl::bigint * 1000)::bigint; |
|||
EXECUTE format( |
|||
'WITH deleted AS (DELETE FROM event WHERE ts < %L::bigint AND (event_type = %L::varchar OR event_type = %L::varchar) RETURNING *) SELECT count(*) FROM deleted', debug_ttl_ts, 'DEBUG_RULE_NODE', 'DEBUG_RULE_CHAIN') into debug_ttl_deleted_count; |
|||
END IF; |
|||
RAISE NOTICE 'Events removed by ttl: %', ttl_deleted_count; |
|||
RAISE NOTICE 'Debug Events removed by ttl: %', debug_ttl_deleted_count; |
|||
deleted := ttl_deleted_count + debug_ttl_deleted_count; |
|||
END |
|||
$$; |
|||
@ -0,0 +1,37 @@ |
|||
/** |
|||
* 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.actors; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
|
|||
import java.util.function.Predicate; |
|||
|
|||
@RequiredArgsConstructor |
|||
public class TbEntityTypeActorIdPredicate implements Predicate<TbActorId> { |
|||
|
|||
private final EntityType entityType; |
|||
|
|||
@Override |
|||
public boolean test(TbActorId actorId) { |
|||
return actorId instanceof TbEntityActorId && testEntityId(((TbEntityActorId) actorId).getEntityId()); |
|||
} |
|||
|
|||
protected boolean testEntityId(EntityId entityId) { |
|||
return entityId.getEntityType().equals(entityType); |
|||
} |
|||
} |
|||
@ -1,37 +0,0 @@ |
|||
/** |
|||
* 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.actors.device; |
|||
|
|||
import akka.actor.ActorRef; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.TbActorMsg; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
|
|||
/** |
|||
* Created by ashvayka on 15.03.18. |
|||
*/ |
|||
@Data |
|||
public final class DeviceActorToRuleEngineMsg implements TbActorMsg { |
|||
|
|||
private final ActorRef callbackRef; |
|||
private final TbMsg tbMsg; |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.DEVICE_ACTOR_TO_RULE_ENGINE_MSG; |
|||
} |
|||
} |
|||
@ -1,83 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import akka.actor.ActorRef; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.ActorService; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
import org.thingsboard.server.service.cluster.rpc.GrpcSession; |
|||
import org.thingsboard.server.service.cluster.rpc.GrpcSessionListener; |
|||
import org.thingsboard.server.service.executors.ClusterRpcCallbackExecutorService; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Slf4j |
|||
public class BasicRpcSessionListener implements GrpcSessionListener { |
|||
|
|||
private final ClusterRpcCallbackExecutorService callbackExecutorService; |
|||
private final ActorService service; |
|||
private final ActorRef manager; |
|||
private final ActorRef self; |
|||
|
|||
BasicRpcSessionListener(ActorSystemContext context, ActorRef manager, ActorRef self) { |
|||
this.service = context.getActorService(); |
|||
this.callbackExecutorService = context.getClusterRpcCallbackExecutor(); |
|||
this.manager = manager; |
|||
this.self = self; |
|||
} |
|||
|
|||
@Override |
|||
public void onConnected(GrpcSession session) { |
|||
log.info("[{}][{}] session started", session.getRemoteServer(), getType(session)); |
|||
if (!session.isClient()) { |
|||
manager.tell(new RpcSessionConnectedMsg(session.getRemoteServer(), session.getSessionId()), self); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onDisconnected(GrpcSession session) { |
|||
log.info("[{}][{}] session closed", session.getRemoteServer(), getType(session)); |
|||
manager.tell(new RpcSessionDisconnectedMsg(session.isClient(), session.getRemoteServer()), self); |
|||
} |
|||
|
|||
@Override |
|||
public void onReceiveClusterGrpcMsg(GrpcSession session, ClusterAPIProtos.ClusterMessage clusterMessage) { |
|||
log.trace("Received session actor msg from [{}][{}]: {}", session.getRemoteServer(), getType(session), clusterMessage); |
|||
callbackExecutorService.execute(() -> { |
|||
try { |
|||
service.onReceivedMsg(session.getRemoteServer(), clusterMessage); |
|||
} catch (Exception e) { |
|||
log.debug("[{}][{}] Failed to process cluster message: {}", session.getRemoteServer(), getType(session), clusterMessage, e); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
@Override |
|||
public void onError(GrpcSession session, Throwable t) { |
|||
log.warn("[{}][{}] session got error -> {}", session.getRemoteServer(), getType(session), t); |
|||
manager.tell(new RpcSessionClosedMsg(session.isClient(), session.getRemoteServer()), self); |
|||
session.close(); |
|||
} |
|||
|
|||
private static String getType(GrpcSession session) { |
|||
return session.isClient() ? "Client" : "Server"; |
|||
} |
|||
|
|||
|
|||
} |
|||
@ -1,27 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcBroadcastMsg { |
|||
private final ClusterAPIProtos.ClusterMessage msg; |
|||
} |
|||
@ -1,230 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import akka.actor.ActorRef; |
|||
import akka.actor.OneForOneStrategy; |
|||
import akka.actor.Props; |
|||
import akka.actor.SupervisorStrategy; |
|||
import akka.event.Logging; |
|||
import akka.event.LoggingAdapter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.ContextAwareActor; |
|||
import org.thingsboard.server.actors.service.ContextBasedCreator; |
|||
import org.thingsboard.server.actors.service.DefaultActorService; |
|||
import org.thingsboard.server.common.msg.TbActorMsg; |
|||
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.common.msg.cluster.ServerType; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
import org.thingsboard.server.service.cluster.discovery.ServerInstance; |
|||
import scala.concurrent.duration.Duration; |
|||
|
|||
import java.util.*; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
public class RpcManagerActor extends ContextAwareActor { |
|||
|
|||
private final Map<ServerAddress, SessionActorInfo> sessionActors; |
|||
private final Map<ServerAddress, Queue<ClusterAPIProtos.ClusterMessage>> pendingMsgs; |
|||
private final ServerAddress instance; |
|||
|
|||
private RpcManagerActor(ActorSystemContext systemContext) { |
|||
super(systemContext); |
|||
this.sessionActors = new HashMap<>(); |
|||
this.pendingMsgs = new HashMap<>(); |
|||
this.instance = systemContext.getDiscoveryService().getCurrentServer().getServerAddress(); |
|||
|
|||
systemContext.getDiscoveryService().getOtherServers().stream() |
|||
.filter(otherServer -> otherServer.getServerAddress().compareTo(instance) > 0) |
|||
.forEach(otherServer -> onCreateSessionRequest( |
|||
new RpcSessionCreateRequestMsg(UUID.randomUUID(), otherServer.getServerAddress(), null))); |
|||
} |
|||
|
|||
@Override |
|||
protected boolean process(TbActorMsg msg) { |
|||
//TODO Move everything here, to work with TbActorMsg
|
|||
return false; |
|||
} |
|||
|
|||
@Override |
|||
public void onReceive(Object msg) { |
|||
if (msg instanceof ClusterAPIProtos.ClusterMessage) { |
|||
onMsg((ClusterAPIProtos.ClusterMessage) msg); |
|||
} else if (msg instanceof RpcBroadcastMsg) { |
|||
onMsg((RpcBroadcastMsg) msg); |
|||
} else if (msg instanceof RpcSessionCreateRequestMsg) { |
|||
onCreateSessionRequest((RpcSessionCreateRequestMsg) msg); |
|||
} else if (msg instanceof RpcSessionConnectedMsg) { |
|||
onSessionConnected((RpcSessionConnectedMsg) msg); |
|||
} else if (msg instanceof RpcSessionDisconnectedMsg) { |
|||
onSessionDisconnected((RpcSessionDisconnectedMsg) msg); |
|||
} else if (msg instanceof RpcSessionClosedMsg) { |
|||
onSessionClosed((RpcSessionClosedMsg) msg); |
|||
} else if (msg instanceof ClusterEventMsg) { |
|||
onClusterEvent((ClusterEventMsg) msg); |
|||
} |
|||
} |
|||
|
|||
private void onMsg(RpcBroadcastMsg msg) { |
|||
log.debug("Forwarding msg to session actors {}", msg); |
|||
sessionActors.keySet().forEach(address -> { |
|||
ClusterAPIProtos.ClusterMessage msgWithServerAddress = msg.getMsg() |
|||
.toBuilder() |
|||
.setServerAddress(ClusterAPIProtos.ServerAddress |
|||
.newBuilder() |
|||
.setHost(address.getHost()) |
|||
.setPort(address.getPort()) |
|||
.build()) |
|||
.build(); |
|||
onMsg(msgWithServerAddress); |
|||
}); |
|||
pendingMsgs.values().forEach(queue -> queue.add(msg.getMsg())); |
|||
} |
|||
|
|||
private void onMsg(ClusterAPIProtos.ClusterMessage msg) { |
|||
if (msg.hasServerAddress()) { |
|||
ServerAddress address = new ServerAddress(msg.getServerAddress().getHost(), msg.getServerAddress().getPort(), ServerType.CORE); |
|||
SessionActorInfo session = sessionActors.get(address); |
|||
if (session != null) { |
|||
log.debug("{} Forwarding msg to session actor: {}", address, msg); |
|||
session.getActor().tell(msg, ActorRef.noSender()); |
|||
} else { |
|||
log.debug("{} Storing msg to pending queue: {}", address, msg); |
|||
Queue<ClusterAPIProtos.ClusterMessage> queue = pendingMsgs.get(address); |
|||
if (queue == null) { |
|||
queue = new LinkedList<>(); |
|||
pendingMsgs.put(new ServerAddress( |
|||
msg.getServerAddress().getHost(), msg.getServerAddress().getPort(), ServerType.CORE), queue); |
|||
} |
|||
queue.add(msg); |
|||
} |
|||
} else { |
|||
log.warn("Cluster msg doesn't have server address [{}]", msg); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void postStop() { |
|||
sessionActors.clear(); |
|||
pendingMsgs.clear(); |
|||
} |
|||
|
|||
private void onClusterEvent(ClusterEventMsg msg) { |
|||
ServerAddress server = msg.getServerAddress(); |
|||
if (server.compareTo(instance) > 0) { |
|||
if (msg.isAdded()) { |
|||
onCreateSessionRequest(new RpcSessionCreateRequestMsg(UUID.randomUUID(), server, null)); |
|||
} else { |
|||
onSessionClose(false, server); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private void onSessionConnected(RpcSessionConnectedMsg msg) { |
|||
register(msg.getRemoteAddress(), msg.getId(), context().sender()); |
|||
} |
|||
|
|||
private void onSessionDisconnected(RpcSessionDisconnectedMsg msg) { |
|||
boolean reconnect = msg.isClient() && isRegistered(msg.getRemoteAddress()); |
|||
onSessionClose(reconnect, msg.getRemoteAddress()); |
|||
} |
|||
|
|||
private void onSessionClosed(RpcSessionClosedMsg msg) { |
|||
boolean reconnect = msg.isClient() && isRegistered(msg.getRemoteAddress()); |
|||
onSessionClose(reconnect, msg.getRemoteAddress()); |
|||
} |
|||
|
|||
private boolean isRegistered(ServerAddress address) { |
|||
for (ServerInstance server : systemContext.getDiscoveryService().getOtherServers()) { |
|||
if (server.getServerAddress().equals(address)) { |
|||
return true; |
|||
} |
|||
} |
|||
return false; |
|||
} |
|||
|
|||
private void onSessionClose(boolean reconnect, ServerAddress remoteAddress) { |
|||
log.info("[{}] session closed. Should reconnect: {}", remoteAddress, reconnect); |
|||
SessionActorInfo sessionRef = sessionActors.get(remoteAddress); |
|||
if (sessionRef != null && context().sender() != null && context().sender().equals(sessionRef.actor)) { |
|||
context().stop(sessionRef.actor); |
|||
sessionActors.remove(remoteAddress); |
|||
pendingMsgs.remove(remoteAddress); |
|||
if (reconnect) { |
|||
onCreateSessionRequest(new RpcSessionCreateRequestMsg(sessionRef.sessionId, remoteAddress, null)); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private void onCreateSessionRequest(RpcSessionCreateRequestMsg msg) { |
|||
if (msg.getRemoteAddress() != null) { |
|||
if (!sessionActors.containsKey(msg.getRemoteAddress())) { |
|||
ActorRef actorRef = createSessionActor(msg); |
|||
register(msg.getRemoteAddress(), msg.getMsgUid(), actorRef); |
|||
} |
|||
} else { |
|||
createSessionActor(msg); |
|||
} |
|||
} |
|||
|
|||
private void register(ServerAddress remoteAddress, UUID uuid, ActorRef sender) { |
|||
sessionActors.put(remoteAddress, new SessionActorInfo(uuid, sender)); |
|||
log.info("[{}][{}] Registering session actor.", remoteAddress, uuid); |
|||
Queue<ClusterAPIProtos.ClusterMessage> data = pendingMsgs.remove(remoteAddress); |
|||
if (data != null) { |
|||
log.info("[{}][{}] Forwarding {} pending messages.", remoteAddress, uuid, data.size()); |
|||
data.forEach(msg -> sender.tell(new RpcSessionTellMsg(msg), ActorRef.noSender())); |
|||
} else { |
|||
log.info("[{}][{}] No pending messages to forward.", remoteAddress, uuid); |
|||
} |
|||
} |
|||
|
|||
private ActorRef createSessionActor(RpcSessionCreateRequestMsg msg) { |
|||
log.info("[{}] Creating session actor.", msg.getMsgUid()); |
|||
ActorRef actor = context().actorOf( |
|||
Props.create(new RpcSessionActor.ActorCreator(systemContext, msg.getMsgUid())) |
|||
.withDispatcher(DefaultActorService.RPC_DISPATCHER_NAME)); |
|||
actor.tell(msg, context().self()); |
|||
return actor; |
|||
} |
|||
|
|||
public static class ActorCreator extends ContextBasedCreator<RpcManagerActor> { |
|||
private static final long serialVersionUID = 1L; |
|||
|
|||
public ActorCreator(ActorSystemContext context) { |
|||
super(context); |
|||
} |
|||
|
|||
@Override |
|||
public RpcManagerActor create() { |
|||
return new RpcManagerActor(context); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public SupervisorStrategy supervisorStrategy() { |
|||
return strategy; |
|||
} |
|||
|
|||
private final SupervisorStrategy strategy = new OneForOneStrategy(3, Duration.create("1 minute"), t -> { |
|||
log.warn("Unknown failure", t); |
|||
return SupervisorStrategy.resume(); |
|||
}); |
|||
} |
|||
@ -1,135 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import io.grpc.ManagedChannel; |
|||
import io.grpc.ManagedChannelBuilder; |
|||
import io.grpc.stub.StreamObserver; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.ContextAwareActor; |
|||
import org.thingsboard.server.actors.service.ContextBasedCreator; |
|||
import org.thingsboard.server.common.msg.TbActorMsg; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
import org.thingsboard.server.gen.cluster.ClusterRpcServiceGrpc; |
|||
import org.thingsboard.server.service.cluster.rpc.GrpcSession; |
|||
import org.thingsboard.server.service.cluster.rpc.GrpcSessionListener; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
import static org.thingsboard.server.gen.cluster.ClusterAPIProtos.MessageType.CONNECT_RPC_MESSAGE; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Slf4j |
|||
public class RpcSessionActor extends ContextAwareActor { |
|||
|
|||
|
|||
private final UUID sessionId; |
|||
private GrpcSession session; |
|||
private GrpcSessionListener listener; |
|||
|
|||
private RpcSessionActor(ActorSystemContext systemContext, UUID sessionId) { |
|||
super(systemContext); |
|||
this.sessionId = sessionId; |
|||
} |
|||
|
|||
@Override |
|||
protected boolean process(TbActorMsg msg) { |
|||
//TODO Move everything here, to work with TbActorMsg
|
|||
return false; |
|||
} |
|||
|
|||
@Override |
|||
public void onReceive(Object msg) { |
|||
if (msg instanceof ClusterAPIProtos.ClusterMessage) { |
|||
tell((ClusterAPIProtos.ClusterMessage) msg); |
|||
} else if (msg instanceof RpcSessionCreateRequestMsg) { |
|||
initSession((RpcSessionCreateRequestMsg) msg); |
|||
} |
|||
} |
|||
|
|||
private void tell(ClusterAPIProtos.ClusterMessage msg) { |
|||
if (session != null) { |
|||
session.sendMsg(msg); |
|||
} else { |
|||
log.trace("Failed to send message due to missing session!"); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void postStop() { |
|||
if (session != null) { |
|||
log.info("Closing session -> {}", session.getRemoteServer()); |
|||
try { |
|||
session.close(); |
|||
} catch (RuntimeException e) { |
|||
log.trace("Failed to close session!", e); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private void initSession(RpcSessionCreateRequestMsg msg) { |
|||
log.info("[{}] Initializing session", context().self()); |
|||
ServerAddress remoteServer = msg.getRemoteAddress(); |
|||
listener = new BasicRpcSessionListener(systemContext, context().parent(), context().self()); |
|||
if (msg.getRemoteAddress() == null) { |
|||
// Server session
|
|||
session = new GrpcSession(listener); |
|||
session.setOutputStream(msg.getResponseObserver()); |
|||
session.initInputStream(); |
|||
session.initOutputStream(); |
|||
systemContext.getRpcService().onSessionCreated(msg.getMsgUid(), session.getInputStream()); |
|||
} else { |
|||
// Client session
|
|||
ManagedChannel channel = ManagedChannelBuilder.forAddress(remoteServer.getHost(), remoteServer.getPort()).usePlaintext().build(); |
|||
session = new GrpcSession(remoteServer, listener, channel); |
|||
session.initInputStream(); |
|||
|
|||
ClusterRpcServiceGrpc.ClusterRpcServiceStub stub = ClusterRpcServiceGrpc.newStub(channel); |
|||
StreamObserver<ClusterAPIProtos.ClusterMessage> outputStream = stub.handleMsgs(session.getInputStream()); |
|||
|
|||
session.setOutputStream(outputStream); |
|||
session.initOutputStream(); |
|||
outputStream.onNext(toConnectMsg()); |
|||
} |
|||
} |
|||
|
|||
public static class ActorCreator extends ContextBasedCreator<RpcSessionActor> { |
|||
private static final long serialVersionUID = 1L; |
|||
|
|||
private final UUID sessionId; |
|||
|
|||
public ActorCreator(ActorSystemContext context, UUID sessionId) { |
|||
super(context); |
|||
this.sessionId = sessionId; |
|||
} |
|||
|
|||
@Override |
|||
public RpcSessionActor create() { |
|||
return new RpcSessionActor(context, sessionId); |
|||
} |
|||
} |
|||
|
|||
private ClusterAPIProtos.ClusterMessage toConnectMsg() { |
|||
ServerAddress instance = systemContext.getDiscoveryService().getCurrentServer().getServerAddress(); |
|||
return ClusterAPIProtos.ClusterMessage.newBuilder().setMessageType(CONNECT_RPC_MESSAGE).setServerAddress( |
|||
ClusterAPIProtos.ServerAddress.newBuilder().setHost(instance.getHost()) |
|||
.setPort(instance.getPort()).build()).build(); |
|||
} |
|||
} |
|||
@ -1,29 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcSessionClosedMsg { |
|||
|
|||
private final boolean client; |
|||
private final ServerAddress remoteAddress; |
|||
} |
|||
@ -1,31 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcSessionConnectedMsg { |
|||
|
|||
private final ServerAddress remoteAddress; |
|||
private final UUID id; |
|||
} |
|||
@ -1,35 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import io.grpc.stub.StreamObserver; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcSessionCreateRequestMsg { |
|||
|
|||
private final UUID msgUid; |
|||
private final ServerAddress remoteAddress; |
|||
private final StreamObserver<ClusterAPIProtos.ClusterMessage> responseObserver; |
|||
|
|||
} |
|||
@ -1,29 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcSessionDisconnectedMsg { |
|||
|
|||
private final boolean client; |
|||
private final ServerAddress remoteAddress; |
|||
} |
|||
@ -1,27 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class RpcSessionTellMsg { |
|||
private final ClusterAPIProtos.ClusterMessage msg; |
|||
} |
|||
@ -1,30 +0,0 @@ |
|||
/** |
|||
* 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.actors.rpc; |
|||
|
|||
import akka.actor.ActorRef; |
|||
import lombok.Data; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Data |
|||
public final class SessionActorInfo { |
|||
protected final UUID sessionId; |
|||
protected final ActorRef actor; |
|||
} |
|||
@ -1,48 +0,0 @@ |
|||
/** |
|||
* 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.actors.ruleChain; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; |
|||
import org.thingsboard.server.common.msg.aware.TenantAwareMsg; |
|||
|
|||
import java.io.Serializable; |
|||
|
|||
/** |
|||
* Created by ashvayka on 19.03.18. |
|||
*/ |
|||
@Data |
|||
final class RemoteToRuleChainTellNextMsg extends RuleNodeToRuleChainTellNextMsg implements TenantAwareMsg, RuleChainAwareMsg { |
|||
|
|||
private static final long serialVersionUID = 2459605482321657447L; |
|||
private final TenantId tenantId; |
|||
private final RuleChainId ruleChainId; |
|||
|
|||
public RemoteToRuleChainTellNextMsg(RuleNodeToRuleChainTellNextMsg original, TenantId tenantId, RuleChainId ruleChainId) { |
|||
super(original.getOriginator(), original.getRelationTypes(), original.getMsg()); |
|||
this.tenantId = tenantId; |
|||
this.ruleChainId = ruleChainId; |
|||
} |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG; |
|||
} |
|||
|
|||
} |
|||
@ -1,87 +0,0 @@ |
|||
/** |
|||
* 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.actors.shared; |
|||
|
|||
import akka.actor.ActorContext; |
|||
import akka.actor.ActorRef; |
|||
import akka.actor.Props; |
|||
import akka.actor.UntypedActor; |
|||
import akka.japi.Creator; |
|||
import com.google.common.collect.BiMap; |
|||
import com.google.common.collect.HashBiMap; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.ContextAwareActor; |
|||
import org.thingsboard.server.common.data.SearchTextBased; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.id.UUIDBased; |
|||
import org.thingsboard.server.common.data.page.PageDataIterable; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
|
|||
/** |
|||
* Created by ashvayka on 15.03.18. |
|||
*/ |
|||
@Slf4j |
|||
public abstract class EntityActorsManager<T extends EntityId, A extends UntypedActor, M extends SearchTextBased<? extends UUIDBased>> { |
|||
|
|||
protected final ActorSystemContext systemContext; |
|||
protected final BiMap<T, ActorRef> actors; |
|||
|
|||
public EntityActorsManager(ActorSystemContext systemContext) { |
|||
this.systemContext = systemContext; |
|||
this.actors = HashBiMap.create(); |
|||
} |
|||
|
|||
protected abstract TenantId getTenantId(); |
|||
|
|||
protected abstract String getDispatcherName(); |
|||
|
|||
protected abstract Creator<A> creator(T entityId); |
|||
|
|||
protected abstract PageDataIterable.FetchFunction<M> getFetchEntitiesFunction(); |
|||
|
|||
public void init(ActorContext context) { |
|||
for (M entity : new PageDataIterable<>(getFetchEntitiesFunction(), ContextAwareActor.ENTITY_PACK_LIMIT)) { |
|||
T entityId = (T) entity.getId(); |
|||
log.debug("[{}|{}] Creating entity actor", entityId.getEntityType(), entityId.getId()); |
|||
//TODO: remove this cast making UUIDBased subclass of EntityId an interface and vice versa.
|
|||
ActorRef actorRef = getOrCreateActor(context, entityId); |
|||
visit(entity, actorRef); |
|||
log.debug("[{}|{}] Entity actor created.", entityId.getEntityType(), entityId.getId()); |
|||
} |
|||
} |
|||
|
|||
public void visit(M entity, ActorRef actorRef) { |
|||
} |
|||
|
|||
public ActorRef getOrCreateActor(ActorContext context, T entityId) { |
|||
return actors.computeIfAbsent(entityId, eId -> |
|||
context.actorOf(Props.create(creator(eId)) |
|||
.withDispatcher(getDispatcherName()), eId.toString())); |
|||
} |
|||
|
|||
public void broadcast(Object msg) { |
|||
actors.values().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); |
|||
} |
|||
|
|||
public void remove(T id) { |
|||
actors.remove(id); |
|||
} |
|||
|
|||
} |
|||
@ -1,59 +0,0 @@ |
|||
/** |
|||
* 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.actors.shared.rulechain; |
|||
|
|||
import akka.actor.ActorRef; |
|||
import akka.japi.Creator; |
|||
import lombok.Getter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.ruleChain.RuleChainActor; |
|||
import org.thingsboard.server.actors.shared.EntityActorsManager; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.dao.rule.RuleChainService; |
|||
|
|||
/** |
|||
* Created by ashvayka on 15.03.18. |
|||
*/ |
|||
@Slf4j |
|||
public abstract class RuleChainManager extends EntityActorsManager<RuleChainId, RuleChainActor, RuleChain> { |
|||
|
|||
protected final RuleChainService service; |
|||
@Getter |
|||
protected RuleChain rootChain; |
|||
@Getter |
|||
protected ActorRef rootChainActor; |
|||
|
|||
public RuleChainManager(ActorSystemContext systemContext) { |
|||
super(systemContext); |
|||
this.service = systemContext.getRuleChainService(); |
|||
} |
|||
|
|||
@Override |
|||
public Creator<RuleChainActor> creator(RuleChainId entityId) { |
|||
return new RuleChainActor.ActorCreator(systemContext, getTenantId(), entityId); |
|||
} |
|||
|
|||
@Override |
|||
public void visit(RuleChain entity, ActorRef actorRef) { |
|||
if (entity != null && entity.isRoot()) { |
|||
rootChain = entity; |
|||
rootChainActor = actorRef; |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -1,48 +0,0 @@ |
|||
/** |
|||
* 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.actors.shared.rulechain; |
|||
|
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.DefaultActorService; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.dao.model.ModelConstants; |
|||
|
|||
import java.util.Collections; |
|||
|
|||
public class SystemRuleChainManager extends RuleChainManager { |
|||
|
|||
public SystemRuleChainManager(ActorSystemContext systemContext) { |
|||
super(systemContext); |
|||
} |
|||
|
|||
@Override |
|||
protected FetchFunction<RuleChain> getFetchEntitiesFunction() { |
|||
return link -> new PageData<>(); |
|||
} |
|||
|
|||
@Override |
|||
protected TenantId getTenantId() { |
|||
return ModelConstants.SYSTEM_TENANT; |
|||
} |
|||
|
|||
@Override |
|||
protected String getDispatcherName() { |
|||
return DefaultActorService.SYSTEM_RULE_DISPATCHER_NAME; |
|||
} |
|||
} |
|||
@ -1,54 +0,0 @@ |
|||
/** |
|||
* 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.actors.shared.rulechain; |
|||
|
|||
import akka.actor.ActorContext; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.service.DefaultActorService; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.common.data.rule.RuleChainType; |
|||
|
|||
public class TenantRuleChainManager extends RuleChainManager { |
|||
|
|||
private final TenantId tenantId; |
|||
|
|||
public TenantRuleChainManager(ActorSystemContext systemContext, TenantId tenantId) { |
|||
super(systemContext); |
|||
this.tenantId = tenantId; |
|||
} |
|||
|
|||
@Override |
|||
public void init(ActorContext context) { |
|||
super.init(context); |
|||
} |
|||
|
|||
@Override |
|||
protected TenantId getTenantId() { |
|||
return tenantId; |
|||
} |
|||
|
|||
@Override |
|||
protected String getDispatcherName() { |
|||
return DefaultActorService.TENANT_RULE_DISPATCHER_NAME; |
|||
} |
|||
|
|||
@Override |
|||
protected FetchFunction<RuleChain> getFetchEntitiesFunction() { |
|||
return link -> service.findTenantRuleChainsByType(tenantId, RuleChainType.SYSTEM, link); |
|||
} |
|||
} |
|||
@ -0,0 +1,54 @@ |
|||
/** |
|||
* 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.controller; |
|||
|
|||
import org.springframework.security.access.prepost.PreAuthorize; |
|||
import org.springframework.web.bind.annotation.RequestMapping; |
|||
import org.springframework.web.bind.annotation.RequestMethod; |
|||
import org.springframework.web.bind.annotation.RequestParam; |
|||
import org.springframework.web.bind.annotation.ResponseBody; |
|||
import org.springframework.web.bind.annotation.RestController; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
import java.util.Arrays; |
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping("/api") |
|||
public class QueueController extends BaseController { |
|||
|
|||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|||
@RequestMapping(value = "/tenant/queues", params = {"serviceType"}, method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public List<String> getTenantQueuesByServiceType(@RequestParam String serviceType) throws ThingsboardException { |
|||
checkParameter("serviceType", serviceType); |
|||
try { |
|||
ServiceType type = ServiceType.valueOf(serviceType); |
|||
switch (type) { |
|||
case TB_RULE_ENGINE: |
|||
return Arrays.asList("Main", "HighPriority", "SequentialByOriginator"); |
|||
default: |
|||
return Collections.emptyList(); |
|||
} |
|||
} catch (Exception e) { |
|||
throw handleException(e); |
|||
} |
|||
} |
|||
} |
|||
@ -1,55 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.springframework.util.Assert; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.common.msg.cluster.ServerType; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
|
|||
import static org.thingsboard.server.utils.MiscUtils.missingProperty; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Service |
|||
@Slf4j |
|||
public class CurrentServerInstanceService implements ServerInstanceService { |
|||
|
|||
@Value("${rpc.bind_host}") |
|||
private String rpcHost; |
|||
@Value("${rpc.bind_port}") |
|||
private Integer rpcPort; |
|||
|
|||
private ServerInstance self; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
Assert.hasLength(rpcHost, missingProperty("rpc.bind_host")); |
|||
Assert.notNull(rpcPort, missingProperty("rpc.bind_port")); |
|||
self = new ServerInstance(new ServerAddress(rpcHost, rpcPort, ServerType.CORE)); |
|||
log.info("Current server instance: [{};{}]", self.getHost(), self.getPort()); |
|||
} |
|||
|
|||
@Override |
|||
public ServerInstance getSelf() { |
|||
return self; |
|||
} |
|||
} |
|||
@ -1,33 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
import java.util.List; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
public interface DiscoveryService { |
|||
|
|||
void publishCurrentServer(); |
|||
|
|||
void unpublishCurrentServer(); |
|||
|
|||
ServerInstance getCurrentServer(); |
|||
|
|||
List<ServerInstance> getOtherServers(); |
|||
|
|||
} |
|||
@ -1,28 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
public interface DiscoveryServiceListener { |
|||
|
|||
void onServerAdded(ServerInstance server); |
|||
|
|||
void onServerUpdated(ServerInstance server); |
|||
|
|||
void onServerRemoved(ServerInstance server); |
|||
} |
|||
@ -1,67 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
|||
import org.springframework.context.annotation.DependsOn; |
|||
import org.springframework.stereotype.Service; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Service |
|||
@ConditionalOnProperty(prefix = "zk", value = "enabled", havingValue = "false", matchIfMissing = true) |
|||
@Slf4j |
|||
@DependsOn("environmentLogService") |
|||
public class DummyDiscoveryService implements DiscoveryService { |
|||
|
|||
@Autowired |
|||
private ServerInstanceService serverInstance; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
log.info("Initializing..."); |
|||
} |
|||
|
|||
@Override |
|||
public void publishCurrentServer() { |
|||
//Do nothing
|
|||
} |
|||
|
|||
@Override |
|||
public void unpublishCurrentServer() { |
|||
//Do nothing
|
|||
} |
|||
|
|||
@Override |
|||
public ServerInstance getCurrentServer() { |
|||
return serverInstance.getSelf(); |
|||
} |
|||
|
|||
@Override |
|||
public List<ServerInstance> getOtherServers() { |
|||
return Collections.emptyList(); |
|||
} |
|||
|
|||
|
|||
} |
|||
@ -1,47 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.Getter; |
|||
import lombok.ToString; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@ToString |
|||
@EqualsAndHashCode(exclude = {"serverInfo", "serverAddress"}) |
|||
public final class ServerInstance implements Comparable<ServerInstance> { |
|||
|
|||
@Getter |
|||
private final String host; |
|||
@Getter |
|||
private final int port; |
|||
@Getter |
|||
private final ServerAddress serverAddress; |
|||
|
|||
public ServerInstance(ServerAddress serverAddress) { |
|||
this.serverAddress = serverAddress; |
|||
this.host = serverAddress.getHost(); |
|||
this.port = serverAddress.getPort(); |
|||
} |
|||
|
|||
@Override |
|||
public int compareTo(ServerInstance o) { |
|||
return this.serverAddress.compareTo(o.serverAddress); |
|||
} |
|||
} |
|||
@ -1,24 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
public interface ServerInstanceService { |
|||
|
|||
ServerInstance getSelf(); |
|||
} |
|||
@ -1,330 +0,0 @@ |
|||
/** |
|||
* 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.cluster.discovery; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.apache.commons.lang3.SerializationException; |
|||
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.ChildData; |
|||
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.ApplicationListener; |
|||
import org.springframework.context.annotation.Lazy; |
|||
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.actors.service.ActorService; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.service.cluster.routing.ClusterRoutingService; |
|||
import org.thingsboard.server.service.state.DeviceStateService; |
|||
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; |
|||
import org.thingsboard.server.utils.MiscUtils; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import javax.annotation.PreDestroy; |
|||
import java.util.List; |
|||
import java.util.NoSuchElementException; |
|||
import java.util.concurrent.ExecutorService; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent.Type.CHILD_REMOVED; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Service |
|||
@ConditionalOnProperty(prefix = "zk", value = "enabled", havingValue = "true", matchIfMissing = false) |
|||
@Slf4j |
|||
public class ZkDiscoveryService implements DiscoveryService, 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; |
|||
|
|||
private String zkNodesDir; |
|||
|
|||
@Autowired |
|||
private ServerInstanceService serverInstance; |
|||
|
|||
@Autowired |
|||
@Lazy |
|||
private TelemetrySubscriptionService tsSubService; |
|||
|
|||
@Autowired |
|||
@Lazy |
|||
private DeviceStateService deviceStateService; |
|||
|
|||
@Autowired |
|||
@Lazy |
|||
private ActorService actorService; |
|||
|
|||
@Autowired |
|||
@Lazy |
|||
private ClusterRoutingService routingService; |
|||
|
|||
private ExecutorService reconnectExecutorService; |
|||
|
|||
private CuratorFramework client; |
|||
private PathChildrenCache cache; |
|||
private String nodePath; |
|||
|
|||
private volatile boolean stopped = true; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
log.info("Initializing..."); |
|||
Assert.hasLength(zkUrl, MiscUtils.missingProperty("zk.url")); |
|||
Assert.notNull(zkRetryInterval, MiscUtils.missingProperty("zk.retry_interval_ms")); |
|||
Assert.notNull(zkConnectionTimeout, MiscUtils.missingProperty("zk.connection_timeout_ms")); |
|||
Assert.notNull(zkSessionTimeout, MiscUtils.missingProperty("zk.session_timeout_ms")); |
|||
|
|||
reconnectExecutorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("zk-discovery")); |
|||
|
|||
log.info("Initializing discovery service using ZK connect string: {}", zkUrl); |
|||
|
|||
zkNodesDir = zkDir + "/nodes"; |
|||
initZkClient(); |
|||
} |
|||
|
|||
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 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"); |
|||
} |
|||
|
|||
@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; |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public 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); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ServerInstance getCurrentServer() { |
|||
return serverInstance.getSelf(); |
|||
} |
|||
|
|||
@Override |
|||
public List<ServerInstance> getOtherServers() { |
|||
return cache.getCurrentData().stream() |
|||
.filter(cd -> !cd.getPath().equals(nodePath)) |
|||
.map(cd -> { |
|||
try { |
|||
return new ServerInstance((ServerAddress) SerializationUtils.deserialize(cd.getData())); |
|||
} catch (NoSuchElementException e) { |
|||
log.error("Failed to decode ZK node", e); |
|||
throw new RuntimeException(e); |
|||
} |
|||
}) |
|||
.collect(Collectors.toList()); |
|||
} |
|||
|
|||
@EventListener(ApplicationReadyEvent.class) |
|||
public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) { |
|||
log.info("Received application ready event. Starting current ZK node."); |
|||
if (stopped) { |
|||
log.debug("Ignoring application ready event. Service is stopped."); |
|||
return; |
|||
} |
|||
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 void childEvent(CuratorFramework curatorFramework, PathChildrenCacheEvent pathChildrenCacheEvent) throws Exception { |
|||
if (stopped) { |
|||
log.debug("Ignoring {}. Service is stopped.", pathChildrenCacheEvent); |
|||
return; |
|||
} |
|||
if (client.getState() != CuratorFrameworkState.STARTED) { |
|||
log.debug("Ignoring {}, ZK client is not started, ZK client state [{}]", pathChildrenCacheEvent, client.getState()); |
|||
return; |
|||
} |
|||
ChildData data = pathChildrenCacheEvent.getData(); |
|||
if (data == null) { |
|||
log.debug("Ignoring {} due to empty child data", pathChildrenCacheEvent); |
|||
return; |
|||
} else if (data.getData() == null) { |
|||
log.debug("Ignoring {} due to empty child's data", pathChildrenCacheEvent); |
|||
return; |
|||
} else if (nodePath != null && nodePath.equals(data.getPath())) { |
|||
if (pathChildrenCacheEvent.getType() == CHILD_REMOVED) { |
|||
log.info("ZK node for current instance is somehow deleted."); |
|||
publishCurrentServer(); |
|||
} |
|||
log.debug("Ignoring event about current server {}", pathChildrenCacheEvent); |
|||
return; |
|||
} |
|||
ServerInstance instance; |
|||
try { |
|||
ServerAddress serverAddress = SerializationUtils.deserialize(data.getData()); |
|||
instance = new ServerInstance(serverAddress); |
|||
} catch (SerializationException e) { |
|||
log.error("Failed to decode server instance for node {}", data.getPath(), e); |
|||
throw e; |
|||
} |
|||
log.info("Processing [{}] event for [{}:{}]", pathChildrenCacheEvent.getType(), instance.getHost(), instance.getPort()); |
|||
switch (pathChildrenCacheEvent.getType()) { |
|||
case CHILD_ADDED: |
|||
routingService.onServerAdded(instance); |
|||
tsSubService.onClusterUpdate(); |
|||
deviceStateService.onClusterUpdate(); |
|||
actorService.onServerAdded(instance); |
|||
break; |
|||
case CHILD_UPDATED: |
|||
routingService.onServerUpdated(instance); |
|||
actorService.onServerUpdated(instance); |
|||
break; |
|||
case CHILD_REMOVED: |
|||
routingService.onServerRemoved(instance); |
|||
tsSubService.onClusterUpdate(); |
|||
deviceStateService.onClusterUpdate(); |
|||
actorService.onServerRemoved(instance); |
|||
break; |
|||
default: |
|||
break; |
|||
} |
|||
} |
|||
} |
|||
@ -1,35 +0,0 @@ |
|||
/** |
|||
* 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.cluster.routing; |
|||
|
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.common.msg.cluster.ServerType; |
|||
import org.thingsboard.server.service.cluster.discovery.DiscoveryServiceListener; |
|||
|
|||
import java.util.Optional; |
|||
import java.util.UUID; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
public interface ClusterRoutingService extends DiscoveryServiceListener { |
|||
|
|||
ServerAddress getCurrentServer(); |
|||
|
|||
Optional<ServerAddress> resolveById(EntityId entityId); |
|||
|
|||
} |
|||
@ -1,153 +0,0 @@ |
|||
/** |
|||
* 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.cluster.routing; |
|||
|
|||
import com.google.common.hash.HashCode; |
|||
import com.google.common.hash.HashFunction; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.springframework.util.Assert; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.common.msg.cluster.ServerType; |
|||
import org.thingsboard.server.service.cluster.discovery.DiscoveryService; |
|||
import org.thingsboard.server.service.cluster.discovery.DiscoveryServiceListener; |
|||
import org.thingsboard.server.service.cluster.discovery.ServerInstance; |
|||
import org.thingsboard.server.utils.MiscUtils; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import java.util.Arrays; |
|||
import java.util.Optional; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ConcurrentNavigableMap; |
|||
import java.util.concurrent.ConcurrentSkipListMap; |
|||
|
|||
/** |
|||
* Cluster service implementation based on consistent hash ring |
|||
*/ |
|||
|
|||
@Service |
|||
@Slf4j |
|||
public class ConsistentClusterRoutingService implements ClusterRoutingService { |
|||
|
|||
@Autowired |
|||
private DiscoveryService discoveryService; |
|||
|
|||
@Value("${cluster.hash_function_name}") |
|||
private String hashFunctionName; |
|||
@Value("${cluster.vitrual_nodes_size}") |
|||
private Integer virtualNodesSize; |
|||
|
|||
private ServerInstance currentServer; |
|||
|
|||
private HashFunction hashFunction; |
|||
|
|||
private ConsistentHashCircle[] circles; |
|||
private ConsistentHashCircle rootCircle; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
log.info("Initializing Cluster routing service!"); |
|||
this.hashFunction = MiscUtils.forName(hashFunctionName); |
|||
this.currentServer = discoveryService.getCurrentServer(); |
|||
this.circles = new ConsistentHashCircle[ServerType.values().length]; |
|||
for (ServerType serverType : ServerType.values()) { |
|||
circles[serverType.ordinal()] = new ConsistentHashCircle(); |
|||
} |
|||
rootCircle = circles[ServerType.CORE.ordinal()]; |
|||
addNode(discoveryService.getCurrentServer()); |
|||
for (ServerInstance instance : discoveryService.getOtherServers()) { |
|||
addNode(instance); |
|||
} |
|||
logCircle(); |
|||
log.info("Cluster routing service initialized!"); |
|||
} |
|||
|
|||
@Override |
|||
public ServerAddress getCurrentServer() { |
|||
return discoveryService.getCurrentServer().getServerAddress(); |
|||
} |
|||
|
|||
@Override |
|||
public Optional<ServerAddress> resolveById(EntityId entityId) { |
|||
return resolveByUuid(rootCircle, entityId.getId()); |
|||
} |
|||
|
|||
private Optional<ServerAddress> resolveByUuid(ConsistentHashCircle circle, UUID uuid) { |
|||
Assert.notNull(uuid); |
|||
if (circle.isEmpty()) { |
|||
return Optional.empty(); |
|||
} |
|||
Long hash = hashFunction.newHasher().putLong(uuid.getMostSignificantBits()) |
|||
.putLong(uuid.getLeastSignificantBits()).hash().asLong(); |
|||
if (!circle.containsKey(hash)) { |
|||
ConcurrentNavigableMap<Long, ServerInstance> tailMap = |
|||
circle.tailMap(hash); |
|||
hash = tailMap.isEmpty() ? |
|||
circle.firstKey() : tailMap.firstKey(); |
|||
} |
|||
ServerInstance result = circle.get(hash); |
|||
if (!currentServer.equals(result)) { |
|||
return Optional.of(result.getServerAddress()); |
|||
} else { |
|||
return Optional.empty(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onServerAdded(ServerInstance server) { |
|||
log.info("On server added event: {}", server); |
|||
addNode(server); |
|||
logCircle(); |
|||
} |
|||
|
|||
@Override |
|||
public void onServerUpdated(ServerInstance server) { |
|||
log.debug("Ignoring server onUpdate event: {}", server); |
|||
} |
|||
|
|||
@Override |
|||
public void onServerRemoved(ServerInstance server) { |
|||
log.info("On server removed event: {}", server); |
|||
removeNode(server); |
|||
logCircle(); |
|||
} |
|||
|
|||
private void addNode(ServerInstance instance) { |
|||
for (int i = 0; i < virtualNodesSize; i++) { |
|||
circles[instance.getServerAddress().getServerType().ordinal()].put(hash(instance, i).asLong(), instance); |
|||
} |
|||
} |
|||
|
|||
private void removeNode(ServerInstance instance) { |
|||
for (int i = 0; i < virtualNodesSize; i++) { |
|||
circles[instance.getServerAddress().getServerType().ordinal()].remove(hash(instance, i).asLong()); |
|||
} |
|||
} |
|||
|
|||
private HashCode hash(ServerInstance instance, int i) { |
|||
return hashFunction.newHasher().putString(instance.getHost(), MiscUtils.UTF8).putInt(instance.getPort()).putInt(i).hash(); |
|||
} |
|||
|
|||
private void logCircle() { |
|||
log.trace("Consistent Hash Circle Start"); |
|||
Arrays.asList(circles).forEach(ConsistentHashCircle::log); |
|||
log.trace("Consistent Hash Circle End"); |
|||
} |
|||
|
|||
} |
|||
@ -1,63 +0,0 @@ |
|||
/** |
|||
* 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.cluster.routing; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.service.cluster.discovery.ServerInstance; |
|||
|
|||
import java.util.concurrent.ConcurrentNavigableMap; |
|||
import java.util.concurrent.ConcurrentSkipListMap; |
|||
|
|||
/** |
|||
* Created by ashvayka on 23.09.18. |
|||
*/ |
|||
@Slf4j |
|||
public class ConsistentHashCircle { |
|||
private final ConcurrentNavigableMap<Long, ServerInstance> circle = |
|||
new ConcurrentSkipListMap<>(); |
|||
|
|||
public void put(long hash, ServerInstance instance) { |
|||
circle.put(hash, instance); |
|||
} |
|||
|
|||
public void remove(long hash) { |
|||
circle.remove(hash); |
|||
} |
|||
|
|||
public boolean isEmpty() { |
|||
return circle.isEmpty(); |
|||
} |
|||
|
|||
public boolean containsKey(Long hash) { |
|||
return circle.containsKey(hash); |
|||
} |
|||
|
|||
public ConcurrentNavigableMap<Long, ServerInstance> tailMap(Long hash) { |
|||
return circle.tailMap(hash); |
|||
} |
|||
|
|||
public Long firstKey() { |
|||
return circle.firstKey(); |
|||
} |
|||
|
|||
public ServerInstance get(Long hash) { |
|||
return circle.get(hash); |
|||
} |
|||
|
|||
public void log() { |
|||
circle.entrySet().forEach((e) -> log.debug("{} -> {}", e.getKey(), e.getValue().getServerAddress())); |
|||
} |
|||
} |
|||
@ -1,161 +0,0 @@ |
|||
/** |
|||
* 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.cluster.rpc; |
|||
|
|||
import com.google.protobuf.ByteString; |
|||
import io.grpc.Server; |
|||
import io.grpc.ServerBuilder; |
|||
import io.grpc.stub.StreamObserver; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.actors.rpc.RpcBroadcastMsg; |
|||
import org.thingsboard.server.actors.rpc.RpcSessionCreateRequestMsg; |
|||
import org.thingsboard.server.common.msg.TbActorMsg; |
|||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|||
import org.thingsboard.server.gen.cluster.ClusterAPIProtos; |
|||
import org.thingsboard.server.gen.cluster.ClusterRpcServiceGrpc; |
|||
import org.thingsboard.server.service.cluster.discovery.ServerInstance; |
|||
import org.thingsboard.server.service.cluster.discovery.ServerInstanceService; |
|||
import org.thingsboard.server.service.encoding.DataDecodingEncodingService; |
|||
|
|||
import javax.annotation.PreDestroy; |
|||
import java.io.IOException; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ArrayBlockingQueue; |
|||
import java.util.concurrent.BlockingQueue; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ConcurrentMap; |
|||
|
|||
/** |
|||
* @author Andrew Shvayka |
|||
*/ |
|||
@Service |
|||
@Slf4j |
|||
public class ClusterGrpcService extends ClusterRpcServiceGrpc.ClusterRpcServiceImplBase implements ClusterRpcService { |
|||
|
|||
@Autowired |
|||
private ServerInstanceService instanceService; |
|||
|
|||
@Autowired |
|||
private DataDecodingEncodingService encodingService; |
|||
|
|||
private RpcMsgListener listener; |
|||
|
|||
private Server server; |
|||
|
|||
private ServerInstance instance; |
|||
|
|||
private ConcurrentMap<UUID, BlockingQueue<StreamObserver<ClusterAPIProtos.ClusterMessage>>> pendingSessionMap = |
|||
new ConcurrentHashMap<>(); |
|||
|
|||
public void init(RpcMsgListener listener) { |
|||
this.listener = listener; |
|||
log.info("Initializing RPC service!"); |
|||
instance = instanceService.getSelf(); |
|||
server = ServerBuilder.forPort(instance.getPort()).addService(this).build(); |
|||
log.info("Going to start RPC server using port: {}", instance.getPort()); |
|||
try { |
|||
server.start(); |
|||
} catch (IOException e) { |
|||
log.error("Failed to start RPC server!", e); |
|||
throw new RuntimeException("Failed to start RPC server!"); |
|||
} |
|||
log.info("RPC service initialized!"); |
|||
} |
|||
|
|||
@Override |
|||
public void onSessionCreated(UUID msgUid, StreamObserver<ClusterAPIProtos.ClusterMessage> inputStream) { |
|||
BlockingQueue<StreamObserver<ClusterAPIProtos.ClusterMessage>> queue = pendingSessionMap.remove(msgUid); |
|||
if (queue != null) { |
|||
try { |
|||
queue.put(inputStream); |
|||
} catch (InterruptedException e) { |
|||
log.warn("Failed to report created session!"); |
|||
Thread.currentThread().interrupt(); |
|||
} |
|||
} else { |
|||
log.warn("Failed to lookup pending session!"); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public StreamObserver<ClusterAPIProtos.ClusterMessage> handleMsgs( |
|||
StreamObserver<ClusterAPIProtos.ClusterMessage> responseObserver) { |
|||
log.info("Processing new session."); |
|||
return createSession(new RpcSessionCreateRequestMsg(UUID.randomUUID(), null, responseObserver)); |
|||
} |
|||
|
|||
|
|||
@PreDestroy |
|||
public void stop() { |
|||
if (server != null) { |
|||
log.info("Going to onStop RPC server"); |
|||
server.shutdownNow(); |
|||
try { |
|||
server.awaitTermination(); |
|||
log.info("RPC server stopped!"); |
|||
} catch (InterruptedException e) { |
|||
log.warn("Failed to onStop RPC server!"); |
|||
Thread.currentThread().interrupt(); |
|||
} |
|||
} |
|||
} |
|||
|
|||
|
|||
@Override |
|||
public void broadcast(RpcBroadcastMsg msg) { |
|||
listener.onBroadcastMsg(msg); |
|||
} |
|||
|
|||
private StreamObserver<ClusterAPIProtos.ClusterMessage> createSession(RpcSessionCreateRequestMsg msg) { |
|||
BlockingQueue<StreamObserver<ClusterAPIProtos.ClusterMessage>> queue = new ArrayBlockingQueue<>(1); |
|||
pendingSessionMap.put(msg.getMsgUid(), queue); |
|||
listener.onRpcSessionCreateRequestMsg(msg); |
|||
try { |
|||
StreamObserver<ClusterAPIProtos.ClusterMessage> observer = queue.take(); |
|||
log.info("Processed new session."); |
|||
return observer; |
|||
} catch (Exception e) { |
|||
log.info("Failed to process session.", e); |
|||
throw new RuntimeException(e); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void tell(ClusterAPIProtos.ClusterMessage message) { |
|||
listener.onSendMsg(message); |
|||
} |
|||
|
|||
@Override |
|||
public void tell(ServerAddress serverAddress, TbActorMsg actorMsg) { |
|||
listener.onSendMsg(encodingService.convertToProtoDataMessage(serverAddress, actorMsg)); |
|||
} |
|||
|
|||
@Override |
|||
public void tell(ServerAddress serverAddress, ClusterAPIProtos.MessageType msgType, byte[] data) { |
|||
ClusterAPIProtos.ClusterMessage msg = ClusterAPIProtos.ClusterMessage |
|||
.newBuilder() |
|||
.setServerAddress(ClusterAPIProtos.ServerAddress |
|||
.newBuilder() |
|||
.setHost(serverAddress.getHost()) |
|||
.setPort(serverAddress.getPort()) |
|||
.build()) |
|||
.setMessageType(msgType) |
|||
.setPayload(ByteString.copyFrom(data)).build(); |
|||
listener.onSendMsg(msg); |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue