Browse Source

Merge branch 'master' into feature/image-resources

pull/9542/head
ViacheslavKlimov 3 years ago
parent
commit
344a7afb47
  1. 1315
      application/src/main/data/json/demo/dashboards/gateways.json
  2. 26
      application/src/main/data/json/tenant/dashboards/gateways.json
  3. 4
      application/src/main/java/org/thingsboard/server/controller/WidgetTypeController.java
  4. 1
      application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java
  5. 9
      application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java
  6. 21
      application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java
  7. 7
      application/src/test/java/org/thingsboard/server/controller/HomePageApiTest.java
  8. 26
      application/src/test/java/org/thingsboard/server/controller/WidgetTypeControllerTest.java
  9. 12
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java
  10. 24
      common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java
  11. 13
      common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java
  12. 42
      common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java
  13. 4
      monitoring/src/main/java/org/thingsboard/monitoring/config/MonitoringTarget.java
  14. 3
      monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringConfig.java
  15. 1
      monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java
  16. 4
      monitoring/src/main/java/org/thingsboard/monitoring/data/Latencies.java
  17. 48
      monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java
  18. 10
      monitoring/src/main/java/org/thingsboard/monitoring/data/notification/HighLatencyNotification.java
  19. 15
      monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java
  20. 89
      monitoring/src/main/java/org/thingsboard/monitoring/service/BaseMonitoringService.java
  21. 26
      monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java
  22. 7
      monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportsMonitoringService.java
  23. 436
      monitoring/src/main/resources/root_rule_chain.json
  24. 16
      monitoring/src/main/resources/tb-monitoring.yml
  25. 6
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ThingsBoardDbInstaller.java
  26. 9
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java

1315
application/src/main/data/json/demo/dashboards/gateways.json

File diff suppressed because it is too large

26
application/src/main/data/json/demo/dashboards/gateway.json → application/src/main/data/json/tenant/dashboards/gateways.json

@ -1,5 +1,5 @@
{
"title": "Gateway",
"title": "ThingsBoard IoT Gateways",
"image": null,
"mobileHide": false,
"mobileOrder": null,
@ -40,7 +40,7 @@
"color": "rgba(0, 0, 0, 0.87)",
"padding": "4px",
"settings": {
"entitiesTitle": "Gateway list",
"entitiesTitle": "Gateways list",
"enableSearch": true,
"enableSelectColumnDisplay": false,
"enableStickyHeader": true,
@ -55,7 +55,7 @@
"defaultSortOrder": "entityName",
"useRowStyleFunction": false
},
"title": "New Entities table",
"title": "Gateways list",
"dropShadow": true,
"enableFullscreen": false,
"titleStyle": {
@ -571,11 +571,11 @@
"padding": "8px",
"settings": {
"useMarkdownTextFunction": true,
"markdownTextFunction": "var blockData = '';\nvar connectorsIndex = ctx.actionsApi.getActionDescriptors('elementClick').findIndex(action=>action.name==\"Connecotrs\");\nvar logsIndex = ctx.actionsApi.getActionDescriptors('elementClick').findIndex(action=>action.name==\"Logs\");\nfunction generateMatHeader(index) {\n if( index !== undefined && index > -1) {\n return `<mat-card-header class='tb-home-widget-link' (click)=\"ctx.actionsApi.handleWidgetAction($event, ctx.actionsApi.getActionDescriptors('elementClick')[${index}], ctx.datasources[0].entity.id)\">`\n } else {\n return \"<mat-card-header >\" \n }\n}\nfunction createDataBlock(value, label, dividerStyle, mobile, index) {\n blockData += `\n <mat-card style=\"flex-grow: 1; width: ${mobile? '100%': 'auto'}; min-height: ${mobile? 'auto': '57px'}\" class=\" ${dividerStyle}\">\n <div class=\"divider\"></div>\n <mat-divider vertical style=\"height:100%\"></mat-divider>\n ${generateMatHeader(index)}\n <mat-card-subtitle>${label}</mat-card-subtitle>\n </mat-card-header>\n <mat-card-content> ${value}</mat-card-content>\n </mat-card>`;\n}\ncreateDataBlock(data[0].Status, \"Status\", data[0].Status === \"Active\"? 'divider-green' : 'divider-red');\ncreateDataBlock(data[0].Name, \"Gateway Name\", '', ctx.isMobile);\ncreateDataBlock(data[0].Type, \"Gateway Type\", '');\ncreateDataBlock(\n `<span style=\"color:rgb(25,128,56)\">${(data[1]?data[1].count:0)} </span>`\n + \" | \" + \n `<span style=\"color:rgb(203,37,48)\">${(data[2]?data[2][\"count 2\"]:0)} </span>`\n , \"Devices <span class='tb-hint' style='padding-left: 0'>(Active | Inactive)</span>\", '');\ncreateDataBlock(\n `<span style=\"color:rgb(25,128,56)\">${(data[0].active_connectors?JSON.parse(data[0].active_connectors).length:0)} </span>`\n + \" | \" + \n `<span style=\"color:rgb(203,37,48)\">${(data[0].inactive_connectors?JSON.parse(data[0].inactive_connectors).length:0)} </span>`\n , \"Connectors <span class='tb-hint' style='padding-left: 0'>(Active | Inactive)</span>\", '', '', connectorsIndex);\ncreateDataBlock(data[0].ALL_ERRORS_COUNT || 0, \"Errors\", (data[0].ALL_ERRORS_COUNT || 0) === 0 ? 'divider-green' : 'divider-red', '', logsIndex);\nreturn `<div fxLayout=\"row wrap\" fxLayoutGap=\"8px\" class=\"cards-container\">${blockData}</div>`;",
"markdownTextFunction": "var blockData = '';\nvar connectorsIndex = ctx.actionsApi.getActionDescriptors('elementClick').findIndex(action=>action.name==\"Connectors\");\nvar logsIndex = ctx.actionsApi.getActionDescriptors('elementClick').findIndex(action=>action.name==\"Logs\");\nfunction generateMatHeader(index) {\n if( index !== undefined && index > -1) {\n return `<mat-card-header class='tb-home-widget-link' (click)=\"ctx.actionsApi.handleWidgetAction($event, ctx.actionsApi.getActionDescriptors('elementClick')[${index}], ctx.datasources[0].entity.id)\">`\n } else {\n return \"<mat-card-header >\" \n }\n}\nfunction createDataBlock(value, label, dividerStyle, mobile, index) {\n blockData += `\n <mat-card style=\"flex-grow: 1; width: ${mobile? '100%': 'auto'}; min-height: ${mobile? 'auto': '57px'}\" class=\" ${dividerStyle}\">\n <div class=\"divider\"></div>\n <mat-divider vertical style=\"height:100%\"></mat-divider>\n ${generateMatHeader(index)}\n <mat-card-subtitle>${label}</mat-card-subtitle>\n </mat-card-header>\n <mat-card-content> ${value}</mat-card-content>\n </mat-card>`;\n}\ncreateDataBlock(data[0].Status, \"Status\", data[0].Status === \"Active\"? 'divider-green' : 'divider-red');\ncreateDataBlock(data[0].Name, \"Gateway Name\", '', ctx.isMobile);\ncreateDataBlock(data[0].Type, \"Gateway Type\", '');\ncreateDataBlock(\n `<span style=\"color:rgb(25,128,56)\">${(data[1]?data[1].count:0)} </span>`\n + \" | \" + \n `<span style=\"color:rgb(203,37,48)\">${(data[2]?data[2][\"count 2\"]:0)} </span>`\n , \"Devices <span class='tb-hint' style='padding-left: 0'>(Active | Inactive)</span>\", '');\ncreateDataBlock(\n `<span style=\"color:rgb(25,128,56)\">${(data[0].active_connectors?JSON.parse(data[0].active_connectors).length:0)} </span>`\n + \" | \" + \n `<span style=\"color:rgb(203,37,48)\">${(data[0].inactive_connectors?JSON.parse(data[0].inactive_connectors).length:0)} </span>`\n , \"Connectors <span class='tb-hint' style='padding-left: 0'>(Active | Inactive)</span>\", '', '', connectorsIndex);\ncreateDataBlock(data[0].ALL_ERRORS_COUNT || 0, \"Errors\", (data[0].ALL_ERRORS_COUNT || 0) === 0 ? 'divider-green' : 'divider-red', '', logsIndex);\nreturn `<div fxLayout=\"row wrap\" fxLayoutGap=\"8px\" class=\"cards-container\">${blockData}</div>`;",
"applyDefaultMarkdownStyle": false,
"markdownCss": ".divider {\n position: absolute;\n width: 3px;\n top: 8px;\n border-radius: 2px;\n bottom: 8px;\n border: 1px solid rgba(31, 70, 144, 1);\n background-color: rgba(31, 70, 144, 1);\n left: 10px;\n}\n.divider-green .divider {\n border: 1px solid rgb(25,128,56);\n background-color: rgb(25,128,56);\n}\n\n.divider-green .mat-mdc-card-content {\n color: rgb(25,128,56);\n}\n\n.divider-red .divider {\n border: 1px solid rgb(203,37,48);\n background-color: rgb(203,37,48);\n}\n\n.divider-red .mat-mdc-card-content {\n color: rgb(203,37,48);\n}\n\n.mdc-card {\n position: relative;\n padding-left: 10px;\n margin-bottom: 1px;\n}\n\n.mat-mdc-card-subtitle {\n font-weight: 400;\n font-size: 12px;\n}\n\n.mat-mdc-card-header {\n padding: 8px 16px 0;\n}\n\n.mat-mdc-card-content:last-child {\n padding-bottom: 8px;\n font-size: 16px;\n}\n\n.cards-container {\n height: calc(100% - 1px);\n justify-content: stretch;\n align-items: center;\n margin-bottom: 1px;\n}\n\n::ng-deep.tb-home-widget-link > div {\n flex-grow: 1;\n cursor: pointer;\n}\n\n .tb-home-widget-link {\n width: 100%;\n }\n\n .tb-home-widget-link:hover::after{\n color: inherit;\n }\n \n .tb-home-widget-link::after{\n content: 'arrow_forward';\n display: inline-block;\n transform: rotate(315deg);\n font-family: 'Material Icons';\n font-weight: normal;\n font-style: normal;\n font-size: 18px;\n color: rgba(0, 0, 0, 0.12);\n vertical-align: bottom;\n margin-left: 6px;\n}"
},
"title": "New Markdown/HTML Card",
"title": "Connectors",
"showTitleIcon": false,
"iconColor": "rgba(0, 0, 0, 0.87)",
"iconSize": "24px",
@ -599,7 +599,7 @@
"actions": {
"elementClick": [
{
"name": "Connecotrs",
"name": "Connectors",
"icon": "more_horiz",
"useShowWidgetActionFunction": null,
"showWidgetActionFunction": "return true;",
@ -680,7 +680,7 @@
"defaultSortOrder": "-createdTime",
"useRowStyleFunction": false
},
"title": "New Alarms table",
"title": "Alarms",
"dropShadow": true,
"enableFullscreen": false,
"titleStyle": {
@ -1074,7 +1074,7 @@
}
]
},
"title": "New RPC remote shell",
"title": "RPC remote shell",
"dropShadow": true,
"enableFullscreen": true,
"widgetStyle": {
@ -1485,7 +1485,7 @@
}
]
},
"title": "New RPC debug terminal",
"title": "RPC debug terminal",
"dropShadow": true,
"enableFullscreen": true,
"widgetStyle": {},
@ -1868,7 +1868,7 @@
"applyDefaultMarkdownStyle": false,
"markdownCss": ".action-buttons-container {\r\n display: flex;\r\n flex-wrap: wrap;\r\n flex-direction: row;\r\n height: 100%;\r\n width: 100%;\r\n align-content: center;\r\n}\r\n\r\nbutton {\r\n flex-grow: 1;\r\n margin: 10px;\r\n min-width: 150px;\r\n height: auto;\r\n}"
},
"title": "New Markdown/HTML Card",
"title": "Service command",
"showTitleIcon": false,
"iconColor": "rgba(0, 0, 0, 0.87)",
"iconSize": "24px",
@ -1980,7 +1980,7 @@
"applyDefaultMarkdownStyle": false,
"markdownCss": ".action-buttons-container {\r\n display: flex;\r\n flex-wrap: wrap;\r\n flex-direction: row;\r\n height: 100%;\r\n width: 100%;\r\n align-content: start;\r\n}\r\n\r\nbutton {\r\n flex-grow: 1;\r\n margin: 10px;\r\n min-width: 150px;\r\n height: auto;\r\n}"
},
"title": "New Markdown/HTML Card",
"title": "General configuration",
"showTitleIcon": false,
"iconColor": "rgba(0, 0, 0, 0.87)",
"iconSize": "24px",
@ -2145,7 +2145,7 @@
"applyDefaultMarkdownStyle": true,
"markdownCss": ".mat-mdc-form-field-subscript-wrapper {\n display: none !important;\n}"
},
"title": "New Markdown/HTML Card",
"title": "Gateway devices",
"showTitleIcon": false,
"iconColor": "rgba(0, 0, 0, 0.87)",
"iconSize": "24px",
@ -4939,7 +4939,7 @@
"applyDefaultMarkdownStyle": false,
"markdownCss": ".action-container {\r\n display: flex;\r\n flex-wrap: wrap;\r\n flex-direction: row;\r\n height: 100%;\r\n width: 100%;\r\n}\r\n\r\nbutton {\r\n flex-grow: 1;\r\n margin: 10px;\r\n min-width: 150px;\r\n height: auto;\r\n}"
},
"title": "New Markdown/HTML Card",
"title": "Gateway commands",
"showTitleIcon": false,
"iconColor": "rgba(0, 0, 0, 0.87)",
"iconSize": "24px",

4
application/src/main/java/org/thingsboard/server/controller/WidgetTypeController.java

@ -226,7 +226,7 @@ public class WidgetTypeController extends AutoCommitController {
@ApiOperation(value = "Get all Widget types for specified Bundle (getBundleWidgetTypes)",
notes = "Returns an array of Widget Type objects that belong to specified Widget Bundle." + WIDGET_TYPE_DESCRIPTION + " " + SYSTEM_OR_TENANT_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/widgetTypes", params = {"widgetsBundleId"}, method = RequestMethod.GET)
@ResponseBody
public List<WidgetType> getBundleWidgetTypes(
@ -258,7 +258,7 @@ public class WidgetTypeController extends AutoCommitController {
@ApiOperation(value = "Get all Widget types details for specified Bundle (getBundleWidgetTypes)",
notes = "Returns an array of Widget Type Details objects that belong to specified Widget Bundle." + WIDGET_TYPE_DETAILS_DESCRIPTION + " " + SYSTEM_OR_TENANT_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/widgetTypesDetails", params = {"widgetsBundleId"}, method = RequestMethod.GET)
@ResponseBody
public List<WidgetTypeDetails> getBundleWidgetTypesDetails(

1
application/src/main/java/org/thingsboard/server/service/entitiy/tenant/DefaultTbTenantService.java

@ -54,6 +54,7 @@ public class DefaultTbTenantService extends AbstractTbEntityService implements T
if (created) {
installScripts.createDefaultRuleChains(savedTenant.getId());
installScripts.createDefaultEdgeRuleChains(savedTenant.getId());
installScripts.createDefaultTenantDashboards(savedTenant.getId(), null);
}
tenantProfileCache.evict(savedTenant.getId());
notificationEntityService.notifyCreateOrUpdateTenant(savedTenant, created ?

9
application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java

@ -347,6 +347,15 @@ public class InstallScripts {
public void loadDashboards(TenantId tenantId, CustomerId customerId) throws Exception {
Path dashboardsDir = Paths.get(getDataDir(), JSON_DIR, DEMO_DIR, DASHBOARDS_DIR);
loadDashboardsFromDir(tenantId, customerId, dashboardsDir);
}
public void createDefaultTenantDashboards(TenantId tenantId, CustomerId customerId) throws Exception {
Path dashboardsDir = Paths.get(getDataDir(), JSON_DIR, TENANT_DIR, DASHBOARDS_DIR);
loadDashboardsFromDir(tenantId, customerId, dashboardsDir);
}
private void loadDashboardsFromDir(TenantId tenantId, CustomerId customerId, Path dashboardsDir) throws IOException {
try (DirectoryStream<Path> dirStream = Files.newDirectoryStream(dashboardsDir, path -> path.toString().endsWith(JSON_EXT))) {
dirStream.forEach(
path -> {

21
application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java

@ -327,7 +327,18 @@ public class DashboardControllerTest extends AbstractControllerTest {
@Test
public void testFindTenantDashboards() throws Exception {
List<DashboardInfo> dashboards = new ArrayList<>();
List<DashboardInfo> expectedDashboards = new ArrayList<>();
PageLink pageLink = new PageLink(24);
PageData<DashboardInfo> pageData = null;
do {
pageData = doGetTypedWithPageLink("/api/tenant/dashboards?",
new TypeReference<PageData<DashboardInfo>>() {
}, pageLink);
expectedDashboards.addAll(pageData.getData());
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} while (pageData.hasNext());
Mockito.reset(tbClusterService, auditLogService);
@ -335,7 +346,7 @@ public class DashboardControllerTest extends AbstractControllerTest {
for (int i = 0; i < cntEntity; i++) {
Dashboard dashboard = new Dashboard();
dashboard.setTitle("Dashboard" + i);
dashboards.add(new DashboardInfo(doPost("/api/dashboard", dashboard, Dashboard.class)));
expectedDashboards.add(new DashboardInfo(doPost("/api/dashboard", dashboard, Dashboard.class)));
}
testNotifyManyEntityManyTimeMsgToEdgeServiceEntityEqAny(new Dashboard(), new Dashboard(),
@ -343,8 +354,6 @@ public class DashboardControllerTest extends AbstractControllerTest {
ActionType.ADDED, cntEntity, cntEntity, cntEntity);
List<DashboardInfo> loadedDashboards = new ArrayList<>();
PageLink pageLink = new PageLink(24);
PageData<DashboardInfo> pageData = null;
do {
pageData = doGetTypedWithPageLink("/api/tenant/dashboards?",
new TypeReference<PageData<DashboardInfo>>() {
@ -355,10 +364,10 @@ public class DashboardControllerTest extends AbstractControllerTest {
}
} while (pageData.hasNext());
dashboards.sort(idComparator);
expectedDashboards.sort(idComparator);
loadedDashboards.sort(idComparator);
Assert.assertEquals(dashboards, loadedDashboards);
Assert.assertEquals(expectedDashboards, loadedDashboards);
}
@Test

7
application/src/test/java/org/thingsboard/server/controller/HomePageApiTest.java

@ -92,6 +92,8 @@ public class HomePageApiTest extends AbstractControllerTest {
@MockBean
private SmsService smsService;
private static final int DEFAULT_DASHBOARDS_COUNT = 1;
//For system administrator
@Test
public void testTenantsCountWsCmd() throws Exception {
@ -408,7 +410,7 @@ public class HomePageApiTest extends AbstractControllerTest {
Assert.assertEquals(2, usageInfo.getUsers());
Assert.assertEquals(configuration.getMaxUsers(), usageInfo.getMaxUsers());
Assert.assertEquals(0, usageInfo.getDashboards());
Assert.assertEquals(DEFAULT_DASHBOARDS_COUNT, usageInfo.getDashboards());
Assert.assertEquals(configuration.getMaxDashboards(), usageInfo.getMaxDashboards());
Assert.assertEquals(0, usageInfo.getTransportMessages());
@ -478,7 +480,8 @@ public class HomePageApiTest extends AbstractControllerTest {
}
usageInfo = doGet("/api/usage", UsageInfo.class);
Assert.assertEquals(dashboards.size(), usageInfo.getDashboards());
int expectedDashboardsCount = dashboards.size() + DEFAULT_DASHBOARDS_COUNT;
Assert.assertEquals(expectedDashboardsCount, usageInfo.getDashboards());
}
private Long getInitialEntityCount(EntityType entityType) throws Exception {

26
application/src/test/java/org/thingsboard/server/controller/WidgetTypeControllerTest.java

@ -190,6 +190,32 @@ public class WidgetTypeControllerTest extends AbstractControllerTest {
Collections.sort(loadedWidgetTypes, idComparator);
Assert.assertEquals(widgetTypes, loadedWidgetTypes);
loginCustomerUser();
List<WidgetType> loadedWidgetTypesCustomer = doGetTyped("/api/widgetTypes?widgetsBundleId={widgetsBundleId}",
new TypeReference<>(){}, widgetsBundle.getId().getId().toString());
Collections.sort(loadedWidgetTypesCustomer, idComparator);
Assert.assertEquals(widgetTypes, loadedWidgetTypesCustomer);
List<WidgetTypeDetails> customerLoadedWidgetTypesDetails = doGetTyped("/api/widgetTypesDetails?widgetsBundleId={widgetsBundleId}",
new TypeReference<>(){}, widgetsBundle.getId().getId().toString());
List<WidgetType> widgetTypesFromDetailsListCustomer = customerLoadedWidgetTypesDetails.stream().map(WidgetType::new).collect(Collectors.toList());
Collections.sort(widgetTypesFromDetailsListCustomer, idComparator);
Assert.assertEquals(widgetTypesFromDetailsListCustomer, loadedWidgetTypes);
loginSysAdmin();
List<WidgetType> sysAdminLoadedWidgetTypes = doGetTyped("/api/widgetTypes?widgetsBundleId={widgetsBundleId}",
new TypeReference<>(){}, widgetsBundle.getId().getId().toString());
Collections.sort(sysAdminLoadedWidgetTypes, idComparator);
Assert.assertEquals(widgetTypes, sysAdminLoadedWidgetTypes);
List<WidgetTypeDetails> sysAdminLoadedWidgetTypesDetails = doGetTyped("/api/widgetTypesDetails?widgetsBundleId={widgetsBundleId}",
new TypeReference<>(){}, widgetsBundle.getId().getId().toString());
List<WidgetType> widgetTypesFromDetailsListSysAdmin = sysAdminLoadedWidgetTypesDetails.stream().map(WidgetType::new).collect(Collectors.toList());
Collections.sort(widgetTypesFromDetailsListSysAdmin, idComparator);
Assert.assertEquals(widgetTypesFromDetailsListSysAdmin, loadedWidgetTypes);
}
@Test

12
common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java

@ -17,6 +17,7 @@ package org.thingsboard.server.cache;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cache.support.NullValue;
import org.springframework.data.redis.connection.RedisConnection;
import org.springframework.data.redis.connection.RedisConnectionFactory;
@ -27,6 +28,7 @@ import org.springframework.data.redis.connection.jedis.JedisConnectionFactory;
import org.springframework.data.redis.core.types.Expiration;
import org.springframework.data.redis.serializer.RedisSerializer;
import org.springframework.data.redis.serializer.StringRedisSerializer;
import org.thingsboard.server.common.data.FstStatsService;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.util.JedisClusterCRC16;
@ -44,6 +46,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
private static final byte[] BINARY_NULL_VALUE = RedisSerializer.java().serialize(NullValue.INSTANCE);
static final JedisPool MOCK_POOL = new JedisPool(); //non-null pool required for JedisConnection to trigger closing jedis connection
@Autowired
private FstStatsService fstStatsService;
@Getter
private final String cacheName;
private final JedisConnectionFactory connectionFactory;
@ -80,6 +85,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return SimpleTbCacheValueWrapper.empty();
} else {
V value = valueSerializer.deserialize(key, rawValue);
if (value != null) {
fstStatsService.incrementDecode(value.getClass());
}
return SimpleTbCacheValueWrapper.wrap(value);
}
}
@ -190,7 +198,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return BINARY_NULL_VALUE;
} else {
try {
return valueSerializer.serialize(value);
var bytes = valueSerializer.serialize(value);
fstStatsService.incrementEncode(value.getClass());
return bytes;
} catch (Exception e) {
log.warn("Failed to serialize the cache value: {}", value, e);
throw new RuntimeException(e);

24
common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.data;
public interface FstStatsService {
void incrementEncode(Class<?> clazz);
void incrementDecode(Class<?> clazz);
}

13
common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java

@ -17,8 +17,10 @@ package org.thingsboard.server.queue.util;
import lombok.extern.slf4j.Slf4j;
import org.nustaq.serialization.FSTConfiguration;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FSTUtils;
import org.thingsboard.server.common.data.FstStatsService;
import java.util.Optional;
@ -26,12 +28,17 @@ import java.util.Optional;
@Service
public class ProtoWithFSTService implements DataDecodingEncodingService {
@Autowired
private FstStatsService fstStatsService;
public static final FSTConfiguration CONFIG = FSTConfiguration.createDefaultConfiguration();
@Override
public <T> Optional<T> decode(byte[] byteArray) {
try {
return Optional.ofNullable(FSTUtils.decode(byteArray));
Optional<T> optional = Optional.ofNullable(FSTUtils.decode(byteArray));
optional.ifPresent(obj -> fstStatsService.incrementDecode(obj.getClass()));
return optional;
} catch (IllegalArgumentException e) {
log.error("Error during deserialization message, [{}]", e.getMessage());
return Optional.empty();
@ -41,7 +48,9 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override
public <T> byte[] encode(T msq) {
return FSTUtils.encode(msq);
var bytes = FSTUtils.encode(msq);
fstStatsService.incrementEncode(msq.getClass());
return bytes;
}

42
common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.stats;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FstStatsService;
import java.util.concurrent.ConcurrentHashMap;
@Service
public class FstStatsServiceImpl implements FstStatsService {
private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>();
@Autowired
private StatsFactory statsFactory;
@Override
public void incrementEncode(Class<?> clazz) {
encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstEncode", key)).increment();
}
@Override
public void incrementDecode(Class<?> clazz) {
decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstDecode", key)).increment();
}
}

4
monitoring/src/main/java/org/thingsboard/monitoring/config/MonitoringTarget.java

@ -21,4 +21,8 @@ public interface MonitoringTarget {
UUID getDeviceId();
String getBaseUrl();
boolean isCheckDomainIps();
}

3
monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringConfig.java

@ -23,9 +23,8 @@ import java.util.List;
@Data
public abstract class TransportMonitoringConfig implements MonitoringConfig<TransportMonitoringTarget> {
private int requestTimeoutMs;
private List<TransportMonitoringTarget> targets;
private int requestTimeoutMs;
public abstract TransportType getTransportType();

1
monitoring/src/main/java/org/thingsboard/monitoring/config/transport/TransportMonitoringTarget.java

@ -25,6 +25,7 @@ public class TransportMonitoringTarget implements MonitoringTarget {
private String baseUrl;
private DeviceConfig device; // set manually during initialization
private boolean checkDomainIps;
@Override
public UUID getDeviceId() {

4
monitoring/src/main/java/org/thingsboard/monitoring/data/Latencies.java

@ -25,4 +25,8 @@ public class Latencies {
return String.format("%sRequest", key);
}
public static String wsUpdate(String key) {
return String.format("%sWsUpdate", key);
}
}

48
monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java

@ -15,53 +15,15 @@
*/
package org.thingsboard.monitoring.data;
import com.google.common.util.concurrent.AtomicDouble;
import lombok.RequiredArgsConstructor;
import lombok.Data;
import java.util.concurrent.atomic.AtomicInteger;
@RequiredArgsConstructor
@Data(staticConstructor = "of")
public class Latency {
private final String key;
private final AtomicDouble latencySum = new AtomicDouble();
private final AtomicInteger counter = new AtomicInteger();
public synchronized void report(double latencyInMs) {
latencySum.addAndGet(latencyInMs);
counter.incrementAndGet();
}
public synchronized double getAvg() {
return latencySum.get() / counter.get();
}
public boolean isNotEmpty() {
return counter.get() > 0;
}
public synchronized void reset() {
latencySum.set(0.0);
counter.set(0);
}
public String getKey() {
return key;
}
public synchronized Latency snapshot() {
Latency snapshot = new Latency(key);
snapshot.latencySum.set(latencySum.get());
snapshot.counter.set(counter.get());
return snapshot;
}
private final double value;
@Override
public String toString() {
return "Latency{" +
"key='" + key + '\'' +
", avgLatency=" + getAvg() +
'}';
public String getFormattedValue() {
return String.format("%,.2f ms", value);
}
}

10
monitoring/src/main/java/org/thingsboard/monitoring/data/notification/HighLatencyNotification.java

@ -21,11 +21,11 @@ import java.util.Collection;
public class HighLatencyNotification implements Notification {
private final Collection<Latency> latencies;
private final Collection<Latency> highLatencies;
private final int thresholdMs;
public HighLatencyNotification(Collection<Latency> latencies, int thresholdMs) {
this.latencies = latencies;
public HighLatencyNotification(Collection<Latency> highLatencies, int thresholdMs) {
this.highLatencies = highLatencies;
this.thresholdMs = thresholdMs;
}
@ -33,8 +33,8 @@ public class HighLatencyNotification implements Notification {
public String getText() {
StringBuilder text = new StringBuilder();
text.append("Some of the latencies are higher than ").append(thresholdMs).append(" ms:\n");
latencies.forEach(latency -> {
text.append(String.format("[%s] %,.2f ms\n", latency.getKey(), latency.getAvg()));
highLatencies.forEach(latency -> {
text.append(String.format("[%s] %s\n", latency.getKey(), latency.getFormattedValue()));
});
return text.toString();
}

15
monitoring/src/main/java/org/thingsboard/monitoring/service/BaseHealthChecker.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.monitoring.service;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -30,13 +31,17 @@ import org.thingsboard.monitoring.util.TbStopWatch;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
@RequiredArgsConstructor
@Slf4j
public abstract class BaseHealthChecker<C extends MonitoringConfig, T extends MonitoringTarget> {
@Getter
protected final C config;
@Getter
protected final T target;
private Object info;
@ -48,6 +53,9 @@ public abstract class BaseHealthChecker<C extends MonitoringConfig, T extends Mo
@Value("${monitoring.check_timeout_ms}")
private int resultCheckTimeoutMs;
@Getter
private final Map<String, BaseHealthChecker<C, T>> associates = new HashMap<>();
public static final String TEST_TELEMETRY_KEY = "testData";
@PostConstruct
@ -84,6 +92,10 @@ public abstract class BaseHealthChecker<C extends MonitoringConfig, T extends Mo
} catch (Exception e) {
reporter.serviceFailure(MonitoredServiceKey.GENERAL, e);
}
associates.values().forEach(healthChecker -> {
healthChecker.check(wsClient);
});
}
private void checkWsUpdate(WsClient wsClient, String testValue) {
@ -96,10 +108,9 @@ public abstract class BaseHealthChecker<C extends MonitoringConfig, T extends Mo
} else if (!update.toString().equals(testValue)) {
throw new ServiceFailureException("Was expecting value " + testValue + " but got " + update);
}
reporter.reportLatency(Latencies.WS_UPDATE, stopWatch.getTime());
reporter.reportLatency(Latencies.wsUpdate(getKey()), stopWatch.getTime());
}
protected abstract void initClient() throws Exception;
protected abstract String createTestPayload(String testValue);

89
monitoring/src/main/java/org/thingsboard/monitoring/service/BaseMonitoringService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.monitoring.service;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationContext;
@ -29,14 +30,22 @@ import org.thingsboard.monitoring.service.transport.TransportHealthChecker;
import org.thingsboard.monitoring.util.TbStopWatch;
import javax.annotation.PostConstruct;
import java.net.InetAddress;
import java.net.URI;
import java.net.URISyntaxException;
import java.util.Arrays;
import java.util.HashSet;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.stream.Collectors;
@Slf4j
public abstract class BaseMonitoringService<C extends MonitoringConfig<T>, T extends MonitoringTarget> {
@Autowired
@Autowired(required = false)
private List<C> configs;
private final List<BaseHealthChecker<C, T>> healthCheckers = new LinkedList<>();
private final List<UUID> devices = new LinkedList<>();
@ -54,18 +63,32 @@ public abstract class BaseMonitoringService<C extends MonitoringConfig<T>, T ext
@PostConstruct
private void init() {
if (configs == null || configs.isEmpty()) {
return;
}
tbClient.logIn();
configs.forEach(config -> {
config.getTargets().forEach(target -> {
BaseHealthChecker<C, T> healthChecker = (BaseHealthChecker<C, T>) createHealthChecker(config, target);
log.info("Initializing {}", healthChecker.getClass().getSimpleName());
healthChecker.initialize(tbClient);
devices.add(target.getDeviceId());
BaseHealthChecker<C, T> healthChecker = initHealthChecker(target, config);
healthCheckers.add(healthChecker);
if (target.isCheckDomainIps()) {
getAssociatedUrls(target.getBaseUrl()).forEach(url -> {
healthChecker.getAssociates().put(url, initHealthChecker(createTarget(url), config));
});
}
});
});
}
private BaseHealthChecker<C, T> initHealthChecker(T target, C config) {
BaseHealthChecker<C, T> healthChecker = (BaseHealthChecker<C, T>) createHealthChecker(config, target);
log.info("Initializing {} for {}", healthChecker.getClass().getSimpleName(), target.getBaseUrl());
healthChecker.initialize(tbClient);
devices.add(target.getDeviceId());
return healthChecker;
}
public final void runChecks() {
if (healthCheckers.isEmpty()) {
return;
@ -78,9 +101,8 @@ public abstract class BaseMonitoringService<C extends MonitoringConfig<T>, T ext
try (WsClient wsClient = wsClientFactory.createClient(accessToken)) {
wsClient.subscribeForTelemetry(devices, TransportHealthChecker.TEST_TELEMETRY_KEY).waitForReply();
for (BaseHealthChecker<C, T> healthChecker : healthCheckers) {
healthChecker.check(wsClient);
check(healthChecker, wsClient);
}
}
reporter.reportLatencies(tbClient);
@ -94,8 +116,61 @@ public abstract class BaseMonitoringService<C extends MonitoringConfig<T>, T ext
}
}
private void check(BaseHealthChecker<C, T> healthChecker, WsClient wsClient) throws Exception {
healthChecker.check(wsClient);
T target = healthChecker.getTarget();
if (target.isCheckDomainIps()) {
Set<String> associatedUrls = getAssociatedUrls(target.getBaseUrl());
Map<String, BaseHealthChecker<C, T>> associates = healthChecker.getAssociates();
Set<String> prevAssociatedUrls = new HashSet<>(associates.keySet());
boolean changed = false;
for (String url : associatedUrls) {
if (!prevAssociatedUrls.contains(url)) {
BaseHealthChecker<C, T> associate = initHealthChecker(createTarget(url), healthChecker.getConfig());
associates.put(url, associate);
changed = true;
}
}
for (String url : prevAssociatedUrls) {
if (!associatedUrls.contains(url)) {
stopHealthChecker(healthChecker);
associates.remove(url);
changed = true;
}
}
if (changed) {
log.info("Updated IPs for {}: {} (old list: {})", target.getBaseUrl(), associatedUrls, prevAssociatedUrls);
}
}
}
@SneakyThrows
private Set<String> getAssociatedUrls(String baseUrl) {
URI url = new URI(baseUrl);
return Arrays.stream(InetAddress.getAllByName(url.getHost()))
.map(InetAddress::getHostAddress)
.map(ip -> {
try {
return new URI(url.getScheme(), null, ip, url.getPort(), "", null, null).toString();
} catch (URISyntaxException e) {
throw new RuntimeException(e);
}
})
.collect(Collectors.toSet());
}
private void stopHealthChecker(BaseHealthChecker<C, T> healthChecker) throws Exception {
healthChecker.destroyClient();
devices.remove(healthChecker.getTarget().getDeviceId());
log.info("Stopped {} for {}", healthChecker.getClass().getSimpleName(), healthChecker.getTarget().getBaseUrl());
}
protected abstract BaseHealthChecker<?, ?> createHealthChecker(C config, T target);
protected abstract T createTarget(String baseUrl);
protected abstract String getName();
}

26
monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java

@ -63,25 +63,20 @@ public class MonitoringReporter {
private String reportingAssetId;
public void reportLatencies(TbClient tbClient) {
List<Latency> latencies = this.latencies.values().stream()
.filter(Latency::isNotEmpty)
.map(latency -> {
Latency snapshot = latency.snapshot();
latency.reset();
return snapshot;
})
.collect(Collectors.toList());
if (latencies.isEmpty()) {
return;
}
log.info("Latencies:\n{}", latencies.stream().map(latency -> latency.getKey() + ": " + latency.getAvg() + " ms")
log.debug("Latencies:\n{}", latencies.values().stream().map(latency -> latency.getKey() + ": " + latency.getFormattedValue())
.collect(Collectors.joining("\n")) + "\n");
if (!latencyReportingEnabled) return;
if (latencies.stream().anyMatch(latency -> latency.getAvg() >= (double) latencyThresholdMs)) {
HighLatencyNotification highLatencyNotification = new HighLatencyNotification(latencies, latencyThresholdMs);
List<Latency> highLatencies = latencies.values().stream()
.filter(latency -> latency.getValue() >= (double) latencyThresholdMs)
.collect(Collectors.toList());
if (!highLatencies.isEmpty()) {
HighLatencyNotification highLatencyNotification = new HighLatencyNotification(highLatencies, latencyThresholdMs);
notificationService.sendNotification(highLatencyNotification);
log.warn("{}", highLatencyNotification.getText());
}
try {
@ -99,10 +94,11 @@ public class MonitoringReporter {
}
ObjectNode msg = JacksonUtil.newObjectNode();
latencies.forEach(latency -> {
msg.set(latency.getKey(), new DoubleNode(latency.getAvg()));
latencies.values().forEach(latency -> {
msg.set(latency.getKey(), new DoubleNode(latency.getValue()));
});
tbClient.saveEntityTelemetry(new AssetId(UUID.fromString(reportingAssetId)), "time", msg);
latencies.clear();
} catch (Exception e) {
log.error("Failed to report latencies: {}", e.getMessage());
}
@ -112,7 +108,7 @@ public class MonitoringReporter {
String latencyKey = key + "Latency";
double latencyInMs = (double) latencyInNanos / 1000_000;
log.trace("Reporting latency [{}]: {} ms", key, latencyInMs);
latencies.computeIfAbsent(latencyKey, k -> new Latency(latencyKey)).report(latencyInMs);
latencies.put(latencyKey, Latency.of(latencyKey, latencyInMs));
}
public void serviceFailure(Object serviceKey, Throwable error) {

7
monitoring/src/main/java/org/thingsboard/monitoring/service/transport/TransportsMonitoringService.java

@ -33,6 +33,13 @@ public final class TransportsMonitoringService extends BaseMonitoringService<Tra
return applicationContext.getBean(config.getTransportType().getServiceClass(), config, target);
}
@Override
protected TransportMonitoringTarget createTarget(String baseUrl) {
TransportMonitoringTarget target = new TransportMonitoringTarget();
target.setBaseUrl(baseUrl);
return target;
}
@Override
protected String getName() {
return "transports check";

436
monitoring/src/main/resources/root_rule_chain.json

@ -0,0 +1,436 @@
{
"ruleChain": {
"additionalInfo": null,
"name": "Root Rule Chain",
"type": "CORE",
"firstRuleNodeId": null,
"root": false,
"debugMode": false,
"configuration": null,
"externalId": null
},
"metadata": {
"firstNodeIndex": 12,
"nodes": [
{
"additionalInfo": {
"description": null,
"layoutX": 1202,
"layoutY": 221
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode",
"name": "Save Timeseries",
"debugMode": true,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"defaultTTL": 0
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 1000,
"layoutY": 167
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "Save Attributes",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 1,
"configuration": {
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": false
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 566,
"layoutY": 302
},
"type": "org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode",
"name": "Message Type Switch",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"version": 0
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 1000,
"layoutY": 381
},
"type": "org.thingsboard.rule.engine.action.TbLogNode",
"name": "Log RPC from Device",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "TBEL",
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);",
"tbelScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);"
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 1000,
"layoutY": 494
},
"type": "org.thingsboard.rule.engine.action.TbLogNode",
"name": "Log Other",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "TBEL",
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);",
"tbelScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);"
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 1000,
"layoutY": 583
},
"type": "org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode",
"name": "RPC Call Request",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"timeoutInSeconds": 60
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 255,
"layoutY": 301
},
"type": "org.thingsboard.rule.engine.filter.TbOriginatorTypeFilterNode",
"name": "Is Entity Group",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"originatorTypes": [
"ENTITY_GROUP"
]
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 319,
"layoutY": 109
},
"type": "org.thingsboard.rule.engine.filter.TbMsgTypeFilterNode",
"name": "Post attributes or RPC request",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"messageTypes": [
"POST_ATTRIBUTES_REQUEST",
"RPC_CALL_FROM_SERVER_TO_DEVICE"
]
},
"externalId": null
},
{
"additionalInfo": {
"layoutX": 627,
"layoutY": 108
},
"type": "org.thingsboard.rule.engine.transform.TbDuplicateMsgToGroupNode",
"name": "Duplicate To Group Entities",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"entityGroupId": null,
"entityGroupIsMessageOriginator": true
},
"externalId": null
},
{
"additionalInfo": {
"description": "Process incoming messages from devices with the alarm rules defined in the device profile. Dispatch all incoming messages with \"Success\" relation type.",
"layoutX": 45,
"layoutY": 359
},
"type": "org.thingsboard.rule.engine.profile.TbDeviceProfileNode",
"name": "Device Profile Node",
"debugMode": true,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"persistAlarmRulesState": false,
"fetchAlarmRulesStateOnStart": false
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 160,
"layoutY": 631
},
"type": "org.thingsboard.rule.engine.filter.TbJsFilterNode",
"name": "Test JS script",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "JS",
"jsScript": "var test = {\n a: 'a',\n b: 'b'\n};\nreturn test.a === 'a' && test.b === 'b';",
"tbelScript": "return msg.temperature > 20;"
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 427,
"layoutY": 541
},
"type": "org.thingsboard.rule.engine.filter.TbJsFilterNode",
"name": "Test TBEL script",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "TBEL",
"jsScript": "return msg.temperature > 20;",
"tbelScript": "var a = \"a\";\nvar b = \"b\";\nreturn a.equals(\"a\") && b.equals(\"b\");"
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 40,
"layoutY": 252
},
"type": "org.thingsboard.rule.engine.transform.TbTransformMsgNode",
"name": "Add arrival timestamp",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "TBEL",
"jsScript": "return {msg: msg, metadata: metadata, msgType: msgType};",
"tbelScript": "metadata.arrivalTs = Date.now();\nreturn {msg: msg, metadata: metadata, msgType: msgType};"
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 1467,
"layoutY": 267
},
"type": "org.thingsboard.rule.engine.transform.TbTransformMsgNode",
"name": "Calculate additional latencies",
"debugMode": true,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"scriptLang": "TBEL",
"jsScript": "return {msg: msg, metadata: metadata, msgType: msgType};",
"tbelScript": "var arrivalLatency = metadata.arrivalTs - metadata.ts;\nvar processingTime = Date.now() - metadata.arrivalTs;\nmsg = {\n arrivalLatency: arrivalLatency,\n processingTime: processingTime\n};\nreturn {msg: msg, metadata: metadata, msgType: msgType};"
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 1438,
"layoutY": 403
},
"type": "org.thingsboard.rule.engine.transform.TbChangeOriginatorNode",
"name": "To latencies asset",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"originatorSource": "ENTITY",
"entityType": "ASSET",
"entityNamePattern": "[Monitoring] Latencies",
"relationsQuery": {
"direction": "FROM",
"maxLevel": 1,
"filters": [
{
"relationType": "Contains",
"entityTypes": []
}
],
"fetchLastLevelOnly": false
}
},
"externalId": null
},
{
"additionalInfo": {
"description": null,
"layoutX": 1458,
"layoutY": 505
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode",
"name": "Save Timeseries",
"debugMode": true,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"defaultTTL": 0
},
"externalId": null
},
{
"additionalInfo": {
"description": "",
"layoutX": 928,
"layoutY": 266
},
"type": "org.thingsboard.rule.engine.filter.TbCheckMessageNode",
"name": "Has testData",
"debugMode": false,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"messageNames": [
"testData"
],
"metadataNames": [],
"checkAllKeys": true
},
"externalId": null
},
{
"additionalInfo": {
"description": null,
"layoutX": 1203,
"layoutY": 327
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode",
"name": "Save Timeseries with TTL",
"debugMode": true,
"singletonMode": false,
"configurationVersion": 0,
"configuration": {
"defaultTTL": 180,
"skipLatestPersistence": null,
"useServerTs": null
},
"externalId": null
}
],
"connections": [
{
"fromIndex": 2,
"toIndex": 1,
"type": "Post attributes"
},
{
"fromIndex": 2,
"toIndex": 3,
"type": "RPC Request from Device"
},
{
"fromIndex": 2,
"toIndex": 4,
"type": "Other"
},
{
"fromIndex": 2,
"toIndex": 5,
"type": "RPC Request to Device"
},
{
"fromIndex": 2,
"toIndex": 16,
"type": "Post telemetry"
},
{
"fromIndex": 6,
"toIndex": 2,
"type": "False"
},
{
"fromIndex": 6,
"toIndex": 7,
"type": "True"
},
{
"fromIndex": 7,
"toIndex": 2,
"type": "False"
},
{
"fromIndex": 7,
"toIndex": 8,
"type": "True"
},
{
"fromIndex": 8,
"toIndex": 2,
"type": "Success"
},
{
"fromIndex": 9,
"toIndex": 10,
"type": "Success"
},
{
"fromIndex": 10,
"toIndex": 11,
"type": "True"
},
{
"fromIndex": 11,
"toIndex": 6,
"type": "True"
},
{
"fromIndex": 12,
"toIndex": 9,
"type": "Success"
},
{
"fromIndex": 13,
"toIndex": 14,
"type": "Success"
},
{
"fromIndex": 14,
"toIndex": 15,
"type": "Success"
},
{
"fromIndex": 16,
"toIndex": 0,
"type": "False"
},
{
"fromIndex": 16,
"toIndex": 17,
"type": "True"
},
{
"fromIndex": 17,
"toIndex": 13,
"type": "Success"
}
],
"ruleChainConnections": null
}
}

16
monitoring/src/main/resources/tb-monitoring.yml

@ -51,8 +51,10 @@ monitoring:
# MQTT QoS
qos: '${MQTT_QOS_LEVEL:1}'
targets:
# MQTT transport base url, tcp://DOMAIN:1883 by default
# MQTT transport base url, tcp://DOMAIN:1883 by default
- base_url: '${MQTT_TRANSPORT_BASE_URL:tcp://${monitoring.domain}:1883}'
# Whether to monitor IPs associated with the domain from base url
check_domain_ips: '${MQTT_TRANSPORT_CHECK_DOMAIN_IPS:false}'
# To add more targets, use following environment variables:
# monitoring.transports.mqtt.targets[1].base_url, monitoring.transports.mqtt.targets[2].base_url, etc.
@ -62,8 +64,10 @@ monitoring:
# CoAP request timeout in milliseconds
request_timeout_ms: '${COAP_REQUEST_TIMEOUT_MS:4000}'
targets:
# CoAP transport base url, coap://DOMAIN by default
# CoAP transport base url, coap://DOMAIN by default
- base_url: '${COAP_TRANSPORT_BASE_URL:coap://${monitoring.domain}}'
# Whether to monitor IPs associated with the domain from base url
check_domain_ips: '${COAP_TRANSPORT_CHECK_DOMAIN_IPS:false}'
# To add more targets, use following environment variables:
# monitoring.transports.coap.targets[1].base_url, monitoring.transports.coap.targets[2].base_url, etc.
@ -73,8 +77,10 @@ monitoring:
# HTTP request timeout in milliseconds
request_timeout_ms: '${HTTP_REQUEST_TIMEOUT_MS:4000}'
targets:
# HTTP transport base url, http://DOMAIN by default
# HTTP transport base url, http://DOMAIN by default
- base_url: '${HTTP_TRANSPORT_BASE_URL:http://${monitoring.domain}}'
# Whether to monitor IPs associated with the domain from base url
check_domain_ips: '${HTTP_TRANSPORT_CHECK_DOMAIN_IPS:false}'
# To add more targets, use following environment variables:
# monitoring.transports.http.targets[1].base_url, monitoring.transports.http.targets[2].base_url, etc.
@ -84,8 +90,10 @@ monitoring:
# LwM2M request timeout in milliseconds
request_timeout_ms: '${LWM2M_REQUEST_TIMEOUT_MS:4000}'
targets:
# LwM2M transport base url, coap://DOMAIN:5685 by default
# LwM2M transport base url, coap://DOMAIN:5685 by default
- base_url: '${LWM2M_TRANSPORT_BASE_URL:coap://${monitoring.domain}:5685}'
# Whether to monitor IPs associated with the domain from base url
check_domain_ips: '${LWM2M_TRANSPORT_CHECK_DOMAIN_IPS:false}'
# To add more targets, use following environment variables:
# monitoring.transports.lwm2m.targets[1].base_url, monitoring.transports.lwm2m.targets[2].base_url, etc.

6
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ThingsBoardDbInstaller.java

@ -253,12 +253,16 @@ public class ThingsBoardDbInstaller {
.add(tbVcExecutorLogVolume)
.add(resolveRedisComposeVolumeLog());
if (IS_HYBRID_MODE) {
rmVolumesCommand.add(cassandraDataVolume);
}
dockerCompose.withCommand(rmVolumesCommand.toString());
}
private String resolveRedisComposeVolumeLog() {
if (IS_REDIS_CLUSTER) {
return IntStream.range(0, 6).mapToObj(i -> redisClusterDataVolume + "-" + i).collect(Collectors.joining());
return IntStream.range(0, 6).mapToObj(i -> " " + redisClusterDataVolume + "-" + i).collect(Collectors.joining());
}
if (IS_REDIS_SENTINEL) {
return redisSentinelDataVolume + "-" + "master " + " " +

9
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java

@ -26,6 +26,7 @@ import io.netty.buffer.Unpooled;
import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
@ -337,8 +338,12 @@ public class MqttClientTest extends AbstractContainerTest {
MqttClient mqttClient = getMqttClient(deviceCredentials, listener);
testRestClient.deleteDeviceIfExists(device.getId());
TimeUnit.SECONDS.sleep(3 * timeoutMultiplier);
assertThat(mqttClient.isConnected()).isFalse();
Awaitility
.await()
.alias("Check device connection.")
.atMost(10, TimeUnit.SECONDS)
.until(() -> !mqttClient.isConnected());
}
@Test

Loading…
Cancel
Save