|
After Width: | Height: | Size: 25 KiB |
|
After Width: | Height: | Size: 15 KiB |
|
After Width: | Height: | Size: 18 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 17 KiB |
|
After Width: | Height: | Size: 16 KiB |
|
After Width: | Height: | Size: 13 KiB |
|
After Width: | Height: | Size: 17 KiB |
|
After Width: | Height: | Size: 30 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 12 KiB |
|
After Width: | Height: | Size: 16 KiB |
|
After Width: | Height: | Size: 15 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 67 KiB |
|
After Width: | Height: | Size: 68 KiB |
|
After Width: | Height: | Size: 11 KiB |
|
After Width: | Height: | Size: 13 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 17 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 17 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 13 KiB |
|
After Width: | Height: | Size: 26 KiB |
|
After Width: | Height: | Size: 14 KiB |
|
After Width: | Height: | Size: 18 KiB |
|
After Width: | Height: | Size: 22 KiB |
|
After Width: | Height: | Size: 15 KiB |
|
After Width: | Height: | Size: 17 KiB |
|
After Width: | Height: | Size: 21 KiB |
|
After Width: | Height: | Size: 19 KiB |
|
After Width: | Height: | Size: 15 KiB |
@ -0,0 +1,48 @@ |
|||
{ |
|||
"widgetsBundle": { |
|||
"alias": "high_performance_scada_energy_system", |
|||
"title": "High-performance SCADA energy system", |
|||
"image": "tb-image:aHBfc2NhZGFfZW5lcmd5X3N5c3RlbV9idW5kbGVfaW1hZ2UucG5n:IkhpZ2gtcGVyZm9ybWFuY2UgU0NBREEgZW5lcmd5IHN5c3RlbSIgc3lzdGVtIGJ1bmRsZSBpbWFnZQ==:SU1BR0U=;data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAMgAAACgCAMAAAB+IdObAAAAllBMVEXe3t7b29vf398AAADf39/e3t7e3t7////Gxsbj4+OpqanFxcW5ubmtra2xsbGVlZWlpaWhoaHS0tLx8fG2trbU1NTBwcGvr6+MjIybm5twcHC+vr6jo6OYmJizs7Orq6vExMSZmZnCwsKdnZ2ampqRkZE3Nze7u7t+fn6EhIRjY2NVVVVsbGx5eXmJiYlEREQvLy8aGhprAy3xAAAABnRSTlPvIK8Av79l/pT7AAAJyElEQVR42uzVXW/aMBiGYWin57VjmyRmNGFLKOWjQACt/P8/t7yJQVMPrE62KkBcR88ZuokTD56HPwa4dU/D58Hw9jNabcYT7sIdHKve4F5C8Ai5No+Qa/MIuTaPkGvzCLk2QSHKIB6rFAIEhayniCc9fVx2c2ou++OUwi88ZGEtYrHbfZOf9/59r9w2f963Fj7hIUaOK0SSWyKSS8s7461rsOmk3TMNj/CQKVGBOJSiTsJNkpjcACjHxOYWHsEhCyJhEYWW1FkDeKNeDcD0c2LgERpi+McrxFFqIYQeaQCqnshWrQAUL7zFGh7BIVNq5YjDJm1Jnfd7lySrnQLT23aX8AgPWbjTG0VFTHKJyohlXLKZ8xy/wiMwxB3rsUEwd0yZ+ecdyQCs+yk1PAJDMuokiMFmxMZl/9W6PJ2EOlrBIzAkp84IUahym6yOG7D0WCarXQE2Oq6Scm3gERoyJxbrA5wJfiCV4l3zljVPu5R8IW7gERZy+fKPC0SgJu4e+XS0Ss+/FSkkI2eKQJ8vxOx7L8RXcjaI4W3GF2KhAKhU8CXYbTHjrZfwCAz5Tc4LohA/pVy7nbZRqds7IVYaHqEhc3ImFgGKkZM1TVO6XR4Oh9rtX+2uRo4BYodYQY7UCCDofwggdoiRdJYDuOEQYuFXYjEqk68qqwKIHZJHCoGQ9FVSgcUNSe8lxNBFBXarIVrSmQG71ZC/1Jpdc6sgEIbv3AXRIB9VE2OStmN6mpnT8/9/3RlSbGTxzmFC9yYTfV14wrIsRLQk/f4ikHDhM+ANVSD6BSB/VLg39SbDJKCyBxFst1o0tkEOsEP2ID20q/nXEFH2IK+Aai1tVVSUOYgrrnYrsx0VEQ2Zg3A6HVq/rlNRkzlI7U/jouMgutsyKmuQEulpHMYHdKW9X8kapIl//hHcIK2Isgbp4zWDRznr9C3KGuQAzkSQorSf60R0zhnEb9G1WF77K7AvYxETm0HYJ0sEUuko2V6lu/Zax6JhM8j+a58IZIjqkav13uuoaum3gtgJJkwDYqJ6ZARvZ0VFm0Nr4oC3NCBv3teBhJGzHS2I240g3MXVjScBEXTzcVm0S0WHbSBsYsAY+6fTgJDJLtOBdB2wzn0kD62yMYdew2z94dSIbaFFB2T/9Ql6YglA5kh6UxUf7X0cwNu9J/b1VCk5izaB7J1nt470t2TpVwup5pGwc/q1P0Q73ymzBYRPrqF7KzdMMEcYAHK79CpHDaDNUcHDEHFraE0ffLaPKUGJYizqwGldnAQKcyyuNrjO0YotIN3+YX2CEbno0OelKCofRgMLb/WiyLZolBYijqKZy5GakbudyhOkbHXM4Yt2UayQgCxzBMER1jiK2oOskZzb/EDqgvqTi7Kx8SJLRcfcQI60oWLwd16W5Qi+RCR5gYQcwYuSPKgiSxkNXE4glEO3D2cYvi6gBNXKfEAoB17oDv1loZaGkuQCQji0sG/0Pyu+BGFIHjjmASKDNrTr+VmRHmB4GoSchTM+B5CWLSisQOKrYvQwnvtiSy9n/PNBSvUYC/7hO4enWswgofNyd+Je33Hxw2KvzwYpR991I6wG9L44ADNVXYYgqhGSwY8I70/5b+fyuSBKAjDRCaPBmX4HZ8a3ik1dDt/OK7VrEVZEWvB3YQFG9VQQ+b+9c1tuGwTCcI+wIJCEhK26suW2aafpRfv+j9ehIgtZuQk20QRn8t9kRP7Y+5nDamUUDfbzr1AzWT2oDTQ65G+YQfiPYPqmBxD3TU0jh9tnHlo74CgGzHLRH5gKqUJ6xGBSTM6mLxwl1NMNrc1FPdLxexzsZq6iFONegCDBdJxPWgBNzvDgwrVnG5GqDavPBemkIBxM+ewHd4Rq/lEtTYREHLvn2njmJnJQB/Nq64sPtcWhFf34AswXjcd/JkvqrGcAoSey3xR+iauiX4sdjrBwTjwG00RIngGk2vJIdcUwRuMNNecCZDTC9pFJn/w0/kfSwjlqk0Eoh5Cdb67ji7o7AA4IEkwu+AFNgpCsKQpCxsQ0YLvDq9FkcdUCrtoKTTEtM1MGSS4Ijoj5KiJDWceFpmiyyz4Eacm2pzrq3WlVEgrSIQeM97deWayiPGzlffFkBm8Kqs3msRm/AgherIaBvOsco41NNk6I+46aUFVtgKNpJVGQnYOo5bB8xw5wZZWzyUfXyhF6nz0ATm986IysnWnHVhIF+dMcW92xpXxhC7EJws0KejoOf5xJeNMJdXq6GRqWqfydDyDIWRvAcgOdciZYcyrkg2hBzqMdmZOKTb0zFQ5CN2doDwJdbLoSkF5HDf2ySDBzzchWUz6IX3JjECzbqKlskNpvckA1WAJSk2HrKR/kK90c5IKmw61N3BxUGQVBzdmTSis4XApiKUjLvczS9JgU+X4rHQE/RMgDORAQMtwkmlJeDCVU6HhtHtbXLhPkBrNfFAztgGPiBwycT7XXPl7nJH9UQw4IjgZJyHzcxJQCgn+lRQQCKRe38kA0TePqBIjhTscUkPoUCKSA6DyQAUEiMjrcRm9KAWmn7SwAsFsvmSA7HjJAfIyGkDntqEmngPzu+cXqf2eA1HRL2ci9FF3JhhQQyTMkM0C2y4TopalpLBrEcqeJBO00UFN7BSBb0kA6AGZT0SDw34TYLkwlg3RAJzbQ/zCAaUCWDOJvBAVCFncAVr9QNIjAcjAiIx2gPEhVMIgWpBxUgl7qx+pXFA1CT7h1jyBVMF0NSK/jBhq38Q26YJA9TeMmAqGmkkFqWg7uA4imJlMwyFdaDtYcpWn1uy8YxD5w20WzMF0ByCFuoHFLNKWCbCRKcREdrAZyQ9O4XILg/SNtMogA1Ib30cFqIIqm8SNHHRamgoeWpuWg4ihLl+RjwSADBTEcJWn1q9LnyGfUrTtArQWCMZq4gcY9XlFCHEMDStEleSgYZEtXqImjDDWNySC9QPXxgVg9j0xx0DRue/byC59Qkm8+Ba0Oso3zCI0b7kwFDy2gSy1w1LQwlQuCFbqKG0gHAKyYEDfKKlB9Jkio0CEOmsQNAk1PDvJ9RpAiF0TgjtWogcStEKR6ahDZI1AeiBakHNQxSEfKeJEAouQZ+o7939+6Y3U5CP26T5N9xdiWBGL5eUKOsDtXZYL0OjTQuA0tfk8LnelC8DAk9WUge1rXNjGIWpoeVqXhLH2KmMBJmeoCkP1u90Pe6dfOCRtCGzG17FE1QogmeYrSpf0CkJZfoG0CCNb7yVUE/1JdN0jV4s3nhYIYlkiigHMx/mQ5IKa5QIYlvTBLVqUbxU7p9Rk9heoVpDS9gpSmV5DS9ApSml7Mg3XffGAvQu9ezkO037+9/sdov/n49v1fKIFpFaO3gooAAAAASUVORK5CYII=", |
|||
"scada": true, |
|||
"description": "Bundle with high-performance SCADA symbols for energy system", |
|||
"order": 9410, |
|||
"name": "High-performance SCADA energy system" |
|||
}, |
|||
"widgetTypeFqns": [ |
|||
"hp_solar_panel", |
|||
"hp_stand_solar_panel", |
|||
"hp_wind_turbine", |
|||
"hp_wind_turbine_cluster", |
|||
"hp_fuel_generator", |
|||
"hp_industrial_fuel_generator", |
|||
"hp_circuit_breaker2", |
|||
"hp_horizontal_circuit_breaker", |
|||
"hp_voltage_relay", |
|||
"hp_3_phase_voltage_relay", |
|||
"hp_voltage_stabilizer", |
|||
"hp_energy_meter", |
|||
"hp_two_rate_energy_meter", |
|||
"hp_three_rate_energy_meter", |
|||
"hp_four_rate_energy_meter", |
|||
"hp_electrical_distribution_board", |
|||
"hp_power_socket", |
|||
"hp_single_key_switch", |
|||
"hp_two_key_switch", |
|||
"hp_top_light_bulb", |
|||
"hp_bottom_light_bulb", |
|||
"hp_battery", |
|||
"hp_inverter", |
|||
"hp_large_inverter", |
|||
"hp_horizontal_energy_systems_controller", |
|||
"hp_vertical_energy_systems_controller", |
|||
"hp_small_power_transformer", |
|||
"hp_power_transformer", |
|||
"hp_low_voltage_transformer_tower", |
|||
"hp_low_voltage_tower", |
|||
"hp_high_voltage_tower", |
|||
"hp_consumers", |
|||
"hp_house", |
|||
"hp_apartments", |
|||
"hp_manufacture" |
|||
] |
|||
} |
|||
@ -1,90 +0,0 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.update; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.rule.engine.api.NotificationCenter; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; |
|||
import org.thingsboard.server.common.data.notification.targets.platform.SystemAdministratorsFilter; |
|||
import org.thingsboard.server.dao.notification.DefaultNotifications; |
|||
import org.thingsboard.server.dao.notification.DefaultNotifications.DefaultNotification; |
|||
import org.thingsboard.server.queue.util.AfterStartUp; |
|||
|
|||
import java.util.Map; |
|||
|
|||
@Service |
|||
@Slf4j |
|||
@RequiredArgsConstructor |
|||
public class DeprecationService { |
|||
|
|||
private final NotificationCenter notificationCenter; |
|||
|
|||
@Value("${queue.type}") |
|||
private String queueType; |
|||
|
|||
@Value("${database.ts.type}") |
|||
private String tsType; |
|||
|
|||
@Value("${database.ts_latest.type}") |
|||
private String tsLatestType; |
|||
|
|||
@AfterStartUp(order = Integer.MAX_VALUE) |
|||
public void checkDeprecation() { |
|||
checkQueueTypeDeprecation(); |
|||
checkDatabaseTypeDeprecation(); |
|||
} |
|||
|
|||
private void checkQueueTypeDeprecation() { |
|||
String queueTypeName; |
|||
switch (queueType) { |
|||
case "aws-sqs" -> queueTypeName = "AWS SQS"; |
|||
case "pubsub" -> queueTypeName = "PubSub"; |
|||
case "service-bus" -> queueTypeName = "Azure Service Bus"; |
|||
case "rabbitmq" -> queueTypeName = "RabbitMQ"; |
|||
default -> { |
|||
return; |
|||
} |
|||
} |
|||
|
|||
log.warn("WARNING: Starting with ThingsBoard 4.0, {} will no longer be supported as a message queue for microservices. " + |
|||
"Please migrate to Apache Kafka. This change will not impact any rule nodes", queueTypeName); |
|||
sendNotification(DefaultNotifications.queueTypeDeprecation, Map.of( |
|||
"queueType", queueTypeName |
|||
)); |
|||
} |
|||
|
|||
private void checkDatabaseTypeDeprecation() { |
|||
String deprecatedDatabaseType = "Timescale"; |
|||
if (StringUtils.equalsAnyIgnoreCase(deprecatedDatabaseType, tsType, tsLatestType)) { |
|||
log.warn("WARNING: Starting with ThingsBoard 4.0, {} will no longer be supported as a storage provider. " + |
|||
"Please migrate to Cassandra or PostgreSQL.", deprecatedDatabaseType); |
|||
sendNotification(DefaultNotifications.databaseTypeDeprecation, Map.of( |
|||
"databaseType", deprecatedDatabaseType |
|||
)); |
|||
} |
|||
} |
|||
|
|||
private void sendNotification(DefaultNotification notification, Map<String, String> info) { |
|||
notificationCenter.sendGeneralWebNotification(TenantId.SYS_TENANT_ID, new SystemAdministratorsFilter(), |
|||
notification.toTemplate(), new GeneralNotificationInfo(info)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,107 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.entitiy.entityview; |
|||
|
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; |
|||
import org.thingsboard.server.common.data.EntityView; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.EntityViewId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|||
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.common.data.objects.AttributesEntityView; |
|||
import org.thingsboard.server.common.data.objects.TelemetryEntityView; |
|||
import org.thingsboard.server.dao.attributes.AttributesService; |
|||
import org.thingsboard.server.dao.entityview.EntityViewService; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; |
|||
|
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
|
|||
import static com.google.common.util.concurrent.Futures.immediateFuture; |
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.mockito.ArgumentMatchers.anyList; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.BDDMockito.given; |
|||
import static org.mockito.BDDMockito.then; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
class DefaultTbEntityViewServiceTest { |
|||
|
|||
final TenantId tenantId = TenantId.fromUUID(UUID.fromString("f09c8180-686c-11ef-9471-a71d33080e9c")); |
|||
final EntityId entityId = DeviceId.fromString("782aaab0-c7a8-11ef-a668-79582e785d5f"); |
|||
|
|||
@Mock |
|||
EntityViewService entityViewService; |
|||
@Mock |
|||
AttributesService attributesService; |
|||
@Mock |
|||
TelemetrySubscriptionService tsSubService; |
|||
@Mock |
|||
TimeseriesService tsService; |
|||
|
|||
DefaultTbEntityViewService defaultTbEntityViewService; |
|||
|
|||
@BeforeEach |
|||
void setup() { |
|||
defaultTbEntityViewService = new DefaultTbEntityViewService(entityViewService, attributesService, tsSubService, tsService); |
|||
} |
|||
|
|||
@Test |
|||
void shouldNotSaveTimeseriesWhenCopyingLatestToEntityView() throws Exception { |
|||
// GIVEN
|
|||
var entityView = new EntityView(new EntityViewId(UUID.randomUUID())); |
|||
entityView.setTenantId(tenantId); |
|||
entityView.setEntityId(entityId); |
|||
entityView.setKeys(new TelemetryEntityView(List.of("temperature"), new AttributesEntityView())); |
|||
|
|||
List<TsKvEntry> latest = List.of(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3))); |
|||
|
|||
given(tsService.findAll(eq(tenantId), eq(entityId), anyList())).willReturn(immediateFuture(latest)); |
|||
|
|||
// WHEN
|
|||
defaultTbEntityViewService.updateEntityViewAttributes(tenantId, entityView, null, null); |
|||
|
|||
// THEN
|
|||
var captor = ArgumentCaptor.forClass(TimeseriesSaveRequest.class); |
|||
then(tsSubService).should().saveTimeseries(captor.capture()); |
|||
|
|||
var expectedCopyLatestRequest = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.entityId(entityView.getId()) |
|||
.entries(latest) |
|||
.ttl(0L) |
|||
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS) |
|||
.build(); |
|||
|
|||
var actualCopyLatestRequest = captor.getValue(); |
|||
|
|||
assertThat(actualCopyLatestRequest) |
|||
.usingRecursiveComparison() |
|||
.ignoringFields("callback") |
|||
.isEqualTo(expectedCopyLatestRequest); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,373 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.telemetry; |
|||
|
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import com.google.common.util.concurrent.SettableFuture; |
|||
import org.checkerframework.checker.nullness.qual.NonNull; |
|||
import org.junit.jupiter.api.AfterEach; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.Arguments; |
|||
import org.junit.jupiter.params.provider.MethodSource; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.springframework.test.util.ReflectionTestUtils; |
|||
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; |
|||
import org.thingsboard.server.cluster.TbClusterService; |
|||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
import org.thingsboard.server.common.data.ApiUsageStateValue; |
|||
import org.thingsboard.server.common.data.EntityView; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.EntityViewId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|||
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|||
import org.thingsboard.server.common.data.kv.KvEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.common.data.objects.AttributesEntityView; |
|||
import org.thingsboard.server.common.data.objects.TelemetryEntityView; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
|||
import org.thingsboard.server.common.stats.TbApiUsageReportClient; |
|||
import org.thingsboard.server.dao.attributes.AttributesService; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.queue.discovery.PartitionService; |
|||
import org.thingsboard.server.queue.discovery.QueueKey; |
|||
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; |
|||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|||
import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; |
|||
import org.thingsboard.server.service.subscription.SubscriptionManagerService; |
|||
|
|||
import java.time.Duration; |
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Optional; |
|||
import java.util.Set; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.ExecutorService; |
|||
import java.util.stream.LongStream; |
|||
import java.util.stream.Stream; |
|||
|
|||
import static com.google.common.util.concurrent.Futures.immediateFuture; |
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.mockito.BDDMockito.given; |
|||
import static org.mockito.BDDMockito.then; |
|||
import static org.mockito.Mockito.lenient; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
class DefaultTelemetrySubscriptionServiceTest { |
|||
|
|||
final TenantId tenantId = TenantId.fromUUID(UUID.fromString("a00ec470-c6b4-11ef-8c88-63b5533fb5bc")); |
|||
final CustomerId customerId = new CustomerId(UUID.fromString("7bdc9750-c775-11ef-8e03-ff69ed8da327")); |
|||
final EntityId entityId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088"); |
|||
|
|||
final long sampleTtl = 10_000L; |
|||
|
|||
final List<TsKvEntry> sampleTelemetry = List.of( |
|||
new BasicTsKvEntry(100L, new DoubleDataEntry("temperature", 65.2)), |
|||
new BasicTsKvEntry(100L, new DoubleDataEntry("humidity", 33.1)) |
|||
); |
|||
|
|||
ApiUsageState apiUsageState; |
|||
|
|||
final TopicPartitionInfo tpi = TopicPartitionInfo.builder() |
|||
.tenantId(tenantId) |
|||
.myPartition(true) |
|||
.build(); |
|||
|
|||
final FutureCallback<Void> emptyCallback = new FutureCallback<>() { |
|||
@Override |
|||
public void onSuccess(Void result) {} |
|||
|
|||
@Override |
|||
public void onFailure(@NonNull Throwable t) {} |
|||
}; |
|||
|
|||
ExecutorService wsCallBackExecutor; |
|||
ExecutorService tsCallBackExecutor; |
|||
|
|||
@Mock |
|||
TbClusterService clusterService; |
|||
@Mock |
|||
PartitionService partitionService; |
|||
@Mock |
|||
SubscriptionManagerService subscriptionManagerService; |
|||
@Mock |
|||
AttributesService attrService; |
|||
@Mock |
|||
TimeseriesService tsService; |
|||
@Mock |
|||
TbEntityViewService tbEntityViewService; |
|||
@Mock |
|||
TbApiUsageReportClient apiUsageClient; |
|||
@Mock |
|||
TbApiUsageStateService apiUsageStateService; |
|||
|
|||
DefaultTelemetrySubscriptionService telemetryService; |
|||
|
|||
@BeforeEach |
|||
void setup() { |
|||
telemetryService = new DefaultTelemetrySubscriptionService(attrService, tsService, tbEntityViewService, apiUsageClient, apiUsageStateService); |
|||
ReflectionTestUtils.setField(telemetryService, "clusterService", clusterService); |
|||
ReflectionTestUtils.setField(telemetryService, "partitionService", partitionService); |
|||
ReflectionTestUtils.setField(telemetryService, "subscriptionManagerService", Optional.of(subscriptionManagerService)); |
|||
|
|||
wsCallBackExecutor = MoreExecutors.newDirectExecutorService(); |
|||
ReflectionTestUtils.setField(telemetryService, "wsCallBackExecutor", wsCallBackExecutor); |
|||
|
|||
tsCallBackExecutor = MoreExecutors.newDirectExecutorService(); |
|||
ReflectionTestUtils.setField(telemetryService, "tsCallBackExecutor", tsCallBackExecutor); |
|||
|
|||
apiUsageState = new ApiUsageState(); |
|||
apiUsageState.setDbStorageState(ApiUsageStateValue.ENABLED); |
|||
lenient().when(apiUsageStateService.getApiUsageState(tenantId)).thenReturn(apiUsageState); |
|||
|
|||
lenient().when(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId)).thenReturn(tpi); |
|||
|
|||
lenient().when(tsService.save(tenantId, entityId, sampleTelemetry, sampleTtl)).thenReturn(immediateFuture(sampleTelemetry.size())); |
|||
lenient().when(tsService.saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl)).thenReturn(immediateFuture(sampleTelemetry.size())); |
|||
lenient().when(tsService.saveLatest(tenantId, entityId, sampleTelemetry)).thenReturn(immediateFuture(listOfNNumbers(sampleTelemetry.size()))); |
|||
|
|||
// mock no entity views
|
|||
lenient().when(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId)).thenReturn(immediateFuture(Collections.emptyList())); |
|||
|
|||
// send partition change event so currentPartitions set is populated
|
|||
telemetryService.onTbApplicationEvent(new PartitionChangeEvent(this, ServiceType.TB_CORE, Map.of(new QueueKey(ServiceType.TB_CORE), Set.of(tpi)), Collections.emptyMap())); |
|||
} |
|||
|
|||
@AfterEach |
|||
void cleanup() { |
|||
wsCallBackExecutor.shutdownNow(); |
|||
tsCallBackExecutor.shutdownNow(); |
|||
} |
|||
|
|||
@Test |
|||
void shouldReportStorageDataPointsApiUsageWhenTimeSeriesIsSaved() { |
|||
// GIVEN
|
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false)) |
|||
.callback(emptyCallback) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
then(apiUsageClient).should().report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, sampleTelemetry.size()); |
|||
} |
|||
|
|||
@Test |
|||
void shouldNotReportStorageDataPointsApiUsageWhenTimeSeriesIsNotSaved() { |
|||
// GIVEN
|
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS) |
|||
.callback(emptyCallback) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
then(apiUsageClient).shouldHaveNoInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
void shouldThrowStorageDisabledWhenTimeSeriesIsSavedAndStorageIsDisabled() { |
|||
// GIVEN
|
|||
apiUsageState.setDbStorageState(ApiUsageStateValue.DISABLED); |
|||
|
|||
SettableFuture<Void> future = SettableFuture.create(); |
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL) |
|||
.future(future) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
assertThat(future).failsWithin(Duration.ofSeconds(5)) |
|||
.withThrowableOfType(ExecutionException.class) |
|||
.withCauseInstanceOf(RuntimeException.class) |
|||
.withMessageContaining("DB storage writes are disabled due to API limits!"); |
|||
} |
|||
|
|||
@Test |
|||
void shouldNotThrowStorageDisabledWhenTimeSeriesIsNotSavedAndStorageIsDisabled() { |
|||
// GIVEN
|
|||
apiUsageState.setDbStorageState(ApiUsageStateValue.DISABLED); |
|||
|
|||
SettableFuture<Void> future = SettableFuture.create(); |
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS) |
|||
.future(future) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
assertThat(future).succeedsWithin(Duration.ofSeconds(5)); |
|||
} |
|||
|
|||
@Test |
|||
void shouldCopyLatestToEntityViewWhenLatestIsSavedOnMainEntity() { |
|||
// GIVEN
|
|||
var entityView = new EntityView(new EntityViewId(UUID.randomUUID())); |
|||
entityView.setTenantId(tenantId); |
|||
entityView.setCustomerId(customerId); |
|||
entityView.setEntityId(entityId); |
|||
entityView.setKeys(new TelemetryEntityView(sampleTelemetry.stream().map(KvEntry::getKey).toList(), new AttributesEntityView())); |
|||
|
|||
// mock that there is one entity view
|
|||
given(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId)).willReturn(immediateFuture(List.of(entityView))); |
|||
// mock that save latest call for entity view is successful
|
|||
given(tsService.saveLatest(tenantId, entityView.getId(), sampleTelemetry)).willReturn(immediateFuture(listOfNNumbers(sampleTelemetry.size()))); |
|||
// mock TPI for entity view
|
|||
given(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityView.getId())).willReturn(tpi); |
|||
|
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(new TimeseriesSaveRequest.Strategy(false, true, false)) |
|||
.callback(emptyCallback) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
// should save latest to both the main entity and it's entity view
|
|||
then(tsService).should().saveLatest(tenantId, entityId, sampleTelemetry); |
|||
then(tsService).should().saveLatest(tenantId, entityView.getId(), sampleTelemetry); |
|||
then(tsService).shouldHaveNoMoreInteractions(); |
|||
|
|||
// should send WS update only for entity view (WS update for the main entity is disabled in the save request)
|
|||
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityView.getId(), sampleTelemetry, TbCallback.EMPTY); |
|||
then(subscriptionManagerService).shouldHaveNoMoreInteractions(); |
|||
} |
|||
|
|||
@Test |
|||
void shouldNotCopyLatestToEntityViewWhenLatestIsNotSavedOnMainEntity() { |
|||
// GIVEN
|
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false)) |
|||
.callback(emptyCallback) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
// should save only time series for the main entity
|
|||
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl); |
|||
then(tsService).shouldHaveNoMoreInteractions(); |
|||
|
|||
// should not send any WS updates
|
|||
then(subscriptionManagerService).shouldHaveNoInteractions(); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@MethodSource("booleanCombinations") |
|||
void shouldCallCorrectApiBasedOnBooleanFlagsInTheSaveRequest(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) { |
|||
// GIVEN
|
|||
var request = TimeseriesSaveRequest.builder() |
|||
.tenantId(tenantId) |
|||
.customerId(customerId) |
|||
.entityId(entityId) |
|||
.entries(sampleTelemetry) |
|||
.ttl(sampleTtl) |
|||
.strategy(new TimeseriesSaveRequest.Strategy(saveTimeseries, saveLatest, sendWsUpdate)) |
|||
.callback(emptyCallback) |
|||
.build(); |
|||
|
|||
// WHEN
|
|||
telemetryService.saveTimeseries(request); |
|||
|
|||
// THEN
|
|||
if (saveTimeseries && saveLatest) { |
|||
then(tsService).should().save(tenantId, entityId, sampleTelemetry, sampleTtl); |
|||
} else if (saveLatest) { |
|||
then(tsService).should().saveLatest(tenantId, entityId, sampleTelemetry); |
|||
} else if (saveTimeseries) { |
|||
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl); |
|||
} |
|||
then(tsService).shouldHaveNoMoreInteractions(); |
|||
|
|||
if (sendWsUpdate) { |
|||
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityId, sampleTelemetry, TbCallback.EMPTY); |
|||
} else { |
|||
then(subscriptionManagerService).shouldHaveNoInteractions(); |
|||
} |
|||
} |
|||
|
|||
private static Stream<Arguments> booleanCombinations() { |
|||
return Stream.of( |
|||
Arguments.of(true, true, true), |
|||
Arguments.of(true, true, false), |
|||
Arguments.of(true, false, true), |
|||
Arguments.of(true, false, false), |
|||
Arguments.of(false, true, true), |
|||
Arguments.of(false, true, false), |
|||
Arguments.of(false, false, true), |
|||
Arguments.of(false, false, false) |
|||
); |
|||
} |
|||
|
|||
// used to emulate sequence numbers returned by save latest API
|
|||
private static List<Long> listOfNNumbers(int N) { |
|||
return LongStream.range(0, N).boxed().toList(); |
|||
} |
|||
|
|||
} |
|||
@ -1,144 +0,0 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.edqs.load; |
|||
|
|||
import lombok.Getter; |
|||
import lombok.RequiredArgsConstructor; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.thingsboard.server.common.data.AttributeScope; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.edqs.AttributeKv; |
|||
import org.thingsboard.server.common.data.edqs.LatestTsKv; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.kv.BooleanDataEntry; |
|||
import org.thingsboard.server.common.data.kv.DataType; |
|||
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|||
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
|||
import org.thingsboard.server.common.data.kv.KvEntry; |
|||
import org.thingsboard.server.common.data.kv.LongDataEntry; |
|||
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|||
import org.thingsboard.server.edqs.processor.EdqsConverter; |
|||
import org.thingsboard.server.edqs.repo.TenantRepo; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
import java.util.Random; |
|||
import java.util.UUID; |
|||
|
|||
@RequiredArgsConstructor |
|||
public class TenantRepoLoader { |
|||
|
|||
private static final int DEVICE_COUNT = 100000; |
|||
private static final int ATTRS_PER_DEVICE = 30; |
|||
private static final int ATTRS_AVG_STR_LENGTH = 12; |
|||
private static final int ATTRS_AVG_JSON_LENGTH = 265; |
|||
private static final int TS_PER_DEVICE = 29; |
|||
private static final int TS_AVG_STR_LENGTH = 59; |
|||
private static final int TS_AVG_JSON_LENGTH = 4005; |
|||
|
|||
private static final Map<DataType, Integer> ATTR_CHANCES = new HashMap<>(); |
|||
private static final Random random = new Random(); |
|||
|
|||
static { |
|||
ATTR_CHANCES.put(DataType.BOOLEAN, 5); |
|||
ATTR_CHANCES.put(DataType.STRING, 49); |
|||
ATTR_CHANCES.put(DataType.LONG, 34); |
|||
ATTR_CHANCES.put(DataType.DOUBLE, 2); |
|||
ATTR_CHANCES.put(DataType.JSON, 10); |
|||
} |
|||
|
|||
private static final Map<DataType, Integer> TS_CHANCES = new HashMap<>(); |
|||
|
|||
static { |
|||
TS_CHANCES.put(DataType.BOOLEAN, 6); |
|||
TS_CHANCES.put(DataType.STRING, 19); |
|||
TS_CHANCES.put(DataType.LONG, 36); |
|||
TS_CHANCES.put(DataType.DOUBLE, 32); |
|||
TS_CHANCES.put(DataType.JSON, 7); |
|||
} |
|||
|
|||
|
|||
@Getter |
|||
private final TenantRepo tenantRepo; |
|||
|
|||
public void load() { |
|||
long ts = System.currentTimeMillis() - DEVICE_COUNT; |
|||
for (int i = 0; i < DEVICE_COUNT; i++) { |
|||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|||
Device device = new Device(); |
|||
device.setId(deviceId); |
|||
device.setCreatedTime(ts + i); |
|||
device.setName("Device " + i); |
|||
device.setLabel("Device Label" + i); |
|||
device.setType("Device Type " + (i % 100)); |
|||
tenantRepo.addOrUpdate(EdqsConverter.toEntity(EntityType.DEVICE, device)); |
|||
for (int j = 0; j < ATTRS_PER_DEVICE; j++) { |
|||
String key = getRandomKey(); |
|||
AttributeKv attributeKv = new AttributeKv(); |
|||
attributeKv.setEntityId(deviceId); |
|||
attributeKv.setScope(AttributeScope.SERVER_SCOPE); |
|||
attributeKv.setKey(key); |
|||
attributeKv.setLastUpdateTs(ts); |
|||
attributeKv.setValue(getRandomKvEntry(key, ATTR_CHANCES, ATTRS_AVG_STR_LENGTH, ATTRS_AVG_JSON_LENGTH)); |
|||
tenantRepo.addOrUpdateAttribute(attributeKv); |
|||
} |
|||
for (int j = 0; j < TS_PER_DEVICE; j++) { |
|||
String key = getRandomKey(); |
|||
LatestTsKv latestTsKv = new LatestTsKv(); |
|||
latestTsKv.setEntityId(deviceId); |
|||
latestTsKv.setKey(key); |
|||
latestTsKv.setTs(ts); |
|||
latestTsKv.setValue(getRandomKvEntry(key, TS_CHANCES, TS_AVG_STR_LENGTH, TS_AVG_JSON_LENGTH)); |
|||
tenantRepo.addOrUpdateLatestKv(latestTsKv); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private KvEntry getRandomKvEntry(String key, Map<DataType, Integer> chances, int strLength, int jsnLength) { |
|||
int i = random.nextInt(100); |
|||
int s = 0; |
|||
for (var pair : chances.entrySet()) { |
|||
s += pair.getValue(); |
|||
if (i < s) { |
|||
switch (pair.getKey()) { |
|||
case BOOLEAN -> { |
|||
return new BooleanDataEntry(key, random.nextBoolean()); |
|||
} |
|||
case LONG -> { |
|||
return new LongDataEntry(key, random.nextLong()); |
|||
} |
|||
case DOUBLE -> { |
|||
return new DoubleDataEntry(key, random.nextDouble()); |
|||
} |
|||
case STRING -> { |
|||
return new StringDataEntry(key, StringUtils.randomAlphanumeric(strLength)); |
|||
} |
|||
case JSON -> { |
|||
return new JsonDataEntry(key, StringUtils.randomAlphanumeric(jsnLength)); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
throw new RuntimeException("Something went wrong"); |
|||
} |
|||
|
|||
private String getRandomKey() { |
|||
return RandomStringUtils.randomAlphabetic(10); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,49 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.transport.mqtt; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.springframework.context.annotation.Bean; |
|||
import org.springframework.context.annotation.Configuration; |
|||
import org.springframework.jmx.export.MBeanExporter; |
|||
import org.thingsboard.server.common.transport.service.DefaultTransportService; |
|||
import org.thingsboard.server.queue.util.TbTransportComponent; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
|
|||
@Configuration |
|||
@TbTransportComponent |
|||
@RequiredArgsConstructor |
|||
public class DefaultTransportMBeanConfiguration { |
|||
|
|||
private final DefaultTransportService transportService; |
|||
|
|||
@Bean |
|||
public HashMapObserver hashMapObserver() { |
|||
return new HashMapObserver(transportService.sessions); |
|||
} |
|||
|
|||
@Bean |
|||
public MBeanExporter mBeanExporter() { |
|||
MBeanExporter exporter = new MBeanExporter(); |
|||
Map<String, Object> beans = new HashMap<>(); |
|||
beans.put("org.thingsboard:type=TransportSessionMapObserver", hashMapObserver()); |
|||
exporter.setBeans(beans); |
|||
return exporter; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,188 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.transport.mqtt; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.common.transport.service.SessionMetaData; |
|||
import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; |
|||
import org.thingsboard.server.transport.mqtt.session.GatewayDeviceSessionContext; |
|||
|
|||
import java.util.Map; |
|||
import java.util.UUID; |
|||
|
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class HashMapObserver implements HashMapObserverMBean { |
|||
private final Map<UUID, SessionMetaData> map; |
|||
|
|||
@Override |
|||
public int getSize() { |
|||
return map.size(); |
|||
} |
|||
|
|||
@Override |
|||
public long getGatewayCount(String unused) { |
|||
return map.values().stream().filter(v-> v.getSessionInfo() != null && v.getSessionInfo().getIsGateway()).count(); |
|||
} |
|||
|
|||
@Override |
|||
public long getNonGatewayCount(String unused) { |
|||
return map.values().stream().filter(v-> v.getSessionInfo() != null && !v.getSessionInfo().getIsGateway()).count(); |
|||
} |
|||
|
|||
@Override |
|||
public String getSessionByUUID(String uuid) { |
|||
return String.valueOf(map.get(UUID.fromString(uuid))); |
|||
} |
|||
|
|||
void addContent(Object entry, int count, StringBuilder content) { |
|||
String lineContent = String.valueOf(entry).replaceAll(System.lineSeparator()," "); |
|||
log.info("{} content = {}", count, lineContent); |
|||
content.append(lineContent).append(System.lineSeparator()); |
|||
} |
|||
|
|||
@Override |
|||
public String getAllSessions(String unused) { |
|||
log.info("getAllSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String getSubscribedSessions(String unused) { |
|||
log.info("getSubscribedSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
boolean hasSubscription = entry.getValue().isSubscribedToRPC() || entry.getValue().isSubscribedToAttributes(); |
|||
if (hasSubscription) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
} |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String getNonActiveSessions(String unused) { |
|||
log.info("getNonActiveSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
SessionMetaData sessionMetaData = entry.getValue(); |
|||
if (sessionMetaData.getListener() instanceof MqttTransportHandler) { |
|||
MqttTransportHandler listener = (MqttTransportHandler) sessionMetaData.getListener(); |
|||
if (!listener.deviceSessionCtx.getChannel().channel().isActive()) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
} |
|||
} |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String getActiveSessions(String unused) { |
|||
log.info("getActiveSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
SessionMetaData sessionMetaData = entry.getValue(); |
|||
if (sessionMetaData.getListener() instanceof MqttTransportHandler) { |
|||
MqttTransportHandler listener = (MqttTransportHandler) sessionMetaData.getListener(); |
|||
if (listener.deviceSessionCtx.getChannel().channel().isActive()) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
} else { |
|||
addContent(sessionMetaData.getListener().getClass(), ++count, content); |
|||
} |
|||
} |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
|
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String getGatewayDeviceSessionContextConnectedSessions(String unused) { |
|||
log.info("getGatewayDeviceSessionContextConnectedSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
SessionMetaData sessionMetaData = entry.getValue(); |
|||
if (sessionMetaData.getListener() instanceof GatewayDeviceSessionContext) { |
|||
GatewayDeviceSessionContext listener = (GatewayDeviceSessionContext) sessionMetaData.getListener(); |
|||
if (listener.isConnected()) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
} |
|||
} |
|||
addContent(count, count, content); |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
throw e; |
|||
} |
|||
} |
|||
@Override |
|||
public String getDeviceAwareSessionContextNotConnectedSessions(String unused) { |
|||
log.info("getDeviceAwareSessionContextNotConnectedSessions()"); |
|||
StringBuilder content = new StringBuilder(); |
|||
try { |
|||
int count = 0; |
|||
for (Map.Entry<UUID, SessionMetaData> entry : map.entrySet()) { |
|||
SessionMetaData sessionMetaData = entry.getValue(); |
|||
if (sessionMetaData.getListener() instanceof DeviceAwareSessionContext) { |
|||
DeviceAwareSessionContext listener = (DeviceAwareSessionContext) sessionMetaData.getListener(); |
|||
if (!listener.isConnected()) { |
|||
addContent(entry, ++count, content); |
|||
} |
|||
} |
|||
} |
|||
addContent(count, count, content); |
|||
return content.toString(); |
|||
} catch (Exception e) { |
|||
log.error(e.getMessage(), e); |
|||
throw e; |
|||
} |
|||
} |
|||
|
|||
} |
|||
|
|||
// 4a7d85c9-eb4b-4fbc-8f6c-deb158cc9ac7=SessionMetaData(sessionInfo=nodeId: "bestia.local" sessionIdMSB: 5367593433178001340 sessionIdLSB: -8111863975520724281 tenantIdMSB: -1954197196874116625 tenantIdLSB: -7022192637061768255 deviceIdMSB: -5222516332438875665 deviceIdLSB: -7338416368958642691 deviceName: "Demo Device" deviceType: "default" gwSessionIdMSB: 3730140660294699909 gwSessionIdLSB: -7918622346767288875 deviceProfileIdMSB: -1952135612572036625 deviceProfileIdLSB: -7022192637061768255 customerIdMSB: 1405474927960789426 customerIdLSB: -9187201950435737472 gatewayIdMSB: 164837549830312431 gatewayIdLSB: -7338416368958642691 , sessionType=ASYNC, listener=GatewayDeviceSessionContext(super=AbstractGatewayDeviceSessionContext(super=MqttDeviceAwareSessionContext(super=DeviceAwareSessionContext(sessionId=4a7d85c9-eb4b-4fbc-8f6c-deb158cc9ac7, deviceId=b785e510-d34e-11ef-9a28-b5316a4ee5fd, tenantId=e4e14c00-d341-11ef-9e8c-29007391dbc1, deviceInfo=TransportDeviceInfo(tenantId=e4e14c00-d341-11ef-9e8c-29007391dbc1, customerId=13814000-1dd2-11b2-8080-808080808080, deviceProfileId=e4e89f00-d341-11ef-9e8c-29007391dbc1, deviceId=b785e510-d34e-11ef-9a28-b5316a4ee5fd, deviceName=Demo Device, deviceType=default, powerMode=null, additionalInfo={"lastConnectedGateway":"02499ee0-d348-11ef-9a28-b5316a4ee5fd"}, edrxCycle=null, psmActivityTimer=null, pagingTransmissionWindow=null, gateway=false), deviceProfile=DeviceProfile(tenantId=e4e14c00-d341-11ef-9e8c-29007391dbc1, name=default, description=Default device profile, isDefault=true, type=DEFAULT, transportType=DEFAULT, provisionType=DISABLED, defaultRuleChainId=null, defaultDashboardId=null, defaultQueueName=null, profileData=DeviceProfileData(configuration=DefaultDeviceProfileConfiguration(), transportConfiguration=DefaultDeviceProfileTransportConfiguration(), provisionConfiguration=DisabledDeviceProfileProvisionConfiguration(provisionDeviceSecret=null), alarms=null), provisionDeviceKey=null, firmwareId=null, softwareId=null, defaultEdgeRuleChainId=null, externalId=null, version=1), sessionInfo=nodeId: "bestia.local" sessionIdMSB: 5367593433178001340 sessionIdLSB: -8111863975520724281 tenantIdMSB: -1954197196874116625 tenantIdLSB: -7022192637061768255 deviceIdMSB: -5222516332438875665 deviceIdLSB: -7338416368958642691 deviceName: "Demo Device" deviceType: "default" gwSessionIdMSB: 3730140660294699909 gwSessionIdLSB: -7918622346767288875 deviceProfileIdMSB: -1952135612572036625 deviceProfileIdLSB: -7022192637061768255 customerIdMSB: 1405474927960789426 customerIdLSB: -9187201950435737472 gatewayIdMSB: 164837549830312431 gatewayIdLSB: -7338416368958642691 , connected=true), mqttQoSMap={org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@697a7658=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@d44c5be8=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@e3a8acd0=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@e3d2681a=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@a65b8616=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@133a8264=1, org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher@12e92606=1}), parent=org.thingsboard.server.transport.mqtt.session.GatewaySessionHandler@3d4885c2, transportService=org.thingsboard.server.common.transport.service.DefaultTransportService@77c16f87)), scheduledFuture=null, subscribedToAttributes=true, subscribedToRPC=true, overwriteActivityTime=false)
|
|||
//1
|
|||
@ -0,0 +1,38 @@ |
|||
/** |
|||
* Copyright © 2016-2024 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.transport.mqtt; |
|||
|
|||
public interface HashMapObserverMBean { |
|||
int getSize(); |
|||
|
|||
long getGatewayCount(String unused); |
|||
|
|||
long getNonGatewayCount(String unused); |
|||
|
|||
String getSessionByUUID(String key); |
|||
|
|||
String getAllSessions(String key); |
|||
|
|||
String getSubscribedSessions(String unused); |
|||
|
|||
String getNonActiveSessions(String unused); |
|||
|
|||
String getActiveSessions(String unused); |
|||
|
|||
String getGatewayDeviceSessionContextConnectedSessions(String unused); |
|||
|
|||
String getDeviceAwareSessionContextNotConnectedSessions(String unused); |
|||
} |
|||