diff --git a/application/pom.xml b/application/pom.xml
index 303f206601..83a9eefdc3 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -289,6 +289,11 @@
junit
test
+
+ org.awaitility
+ awaitility
+ test
+
org.mockito
mockito-core
diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
index 8b240cffed..c3aa3cd9d6 100644
--- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
@@ -381,11 +381,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
}
private void reportSessionOpen() {
- systemContext.getDeviceStateService().onDeviceConnect(deviceId);
+ systemContext.getDeviceStateService().onDeviceConnect(tenantId, deviceId);
}
private void reportSessionClose() {
- systemContext.getDeviceStateService().onDeviceDisconnect(deviceId);
+ systemContext.getDeviceStateService().onDeviceDisconnect(tenantId, deviceId);
}
private void handleGetAttributesRequest(TbActorCtx context, SessionInfoProto sessionInfo, GetAttributeRequestMsg request) {
@@ -590,7 +590,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (sessions.size() == 1) {
reportSessionOpen();
}
- systemContext.getDeviceStateService().onDeviceActivity(deviceId, System.currentTimeMillis());
+ systemContext.getDeviceStateService().onDeviceActivity(tenantId, deviceId, System.currentTimeMillis());
dumpSessions();
} else if (msg.getEvent() == SessionEvent.CLOSED) {
log.debug("[{}] Canceling subscriptions for closed session [{}]", deviceId, sessionId);
@@ -620,7 +620,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (subscriptionInfo.getRpcSubscription()) {
rpcSubscriptions.putIfAbsent(sessionId, sessionMD.getSessionInfo());
}
- systemContext.getDeviceStateService().onDeviceActivity(deviceId, subscriptionInfo.getLastActivityTime());
+ systemContext.getDeviceStateService().onDeviceActivity(tenantId, deviceId, subscriptionInfo.getLastActivityTime());
dumpSessions();
}
diff --git a/application/src/main/java/org/thingsboard/server/controller/AdminController.java b/application/src/main/java/org/thingsboard/server/controller/AdminController.java
index 5bd099d07a..d8406088c7 100644
--- a/application/src/main/java/org/thingsboard/server/controller/AdminController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/AdminController.java
@@ -26,12 +26,12 @@ import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.RestController;
import org.thingsboard.rule.engine.api.MailService;
import org.thingsboard.rule.engine.api.SmsService;
-import org.thingsboard.server.common.data.sms.config.TestSmsRequest;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.UpdateMessage;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
+import org.thingsboard.server.common.data.sms.config.TestSmsRequest;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.permission.Operation;
@@ -67,7 +67,7 @@ public class AdminController extends BaseController {
accessControlService.checkPermission(getCurrentUser(), Resource.ADMIN_SETTINGS, Operation.READ);
AdminSettings adminSettings = checkNotNull(adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, key));
if (adminSettings.getKey().equals("mail")) {
- ((ObjectNode) adminSettings.getJsonValue()).put("password", "");
+ ((ObjectNode) adminSettings.getJsonValue()).remove("password");
}
return adminSettings;
} catch (Exception e) {
@@ -84,7 +84,7 @@ public class AdminController extends BaseController {
adminSettings = checkNotNull(adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, adminSettings));
if (adminSettings.getKey().equals("mail")) {
mailService.updateMailConfiguration();
- ((ObjectNode) adminSettings.getJsonValue()).put("password", "");
+ ((ObjectNode) adminSettings.getJsonValue()).remove("password");
} else if (adminSettings.getKey().equals("sms")) {
smsService.updateSmsConfiguration();
}
@@ -126,6 +126,10 @@ public class AdminController extends BaseController {
accessControlService.checkPermission(getCurrentUser(), Resource.ADMIN_SETTINGS, Operation.READ);
adminSettings = checkNotNull(adminSettings);
if (adminSettings.getKey().equals("mail")) {
+ if(!adminSettings.getJsonValue().has("password")) {
+ AdminSettings mailSettings = checkNotNull(adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "mail"));
+ ((ObjectNode) adminSettings.getJsonValue()).put("password", mailSettings.getJsonValue().get("password").asText());
+ }
String email = getCurrentUser().getEmail();
mailService.sendTestMail(adminSettings.getJsonValue(), email);
}
diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
index 51f645a2af..7042bb6de1 100644
--- a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
+++ b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
@@ -215,6 +215,7 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
node.put("password", "");
node.put("tlsVersion", "TLSv1.2");//NOSONAR, key used to identify password field (not password value itself)
node.put("enableProxy", false);
+ node.put("showChangePassword", false);
mailSettings.setJsonValue(node);
adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, mailSettings);
}
diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CustomOAuth2ClientMapper.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CustomOAuth2ClientMapper.java
index 778f7416ff..93b3fb11b8 100644
--- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CustomOAuth2ClientMapper.java
+++ b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/CustomOAuth2ClientMapper.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.service.security.auth.oauth2;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.web.client.RestTemplateBuilder;
import org.springframework.security.oauth2.client.authentication.OAuth2AuthenticationToken;
@@ -29,6 +30,7 @@ import org.thingsboard.server.common.data.oauth2.OAuth2Registration;
import org.thingsboard.server.dao.oauth2.OAuth2User;
import org.thingsboard.server.service.security.model.SecurityUser;
+import javax.annotation.PostConstruct;
import javax.servlet.http.HttpServletRequest;
@Service(value = "customOAuth2ClientMapper")
@@ -40,6 +42,15 @@ public class CustomOAuth2ClientMapper extends AbstractOAuth2ClientMapper impleme
private RestTemplateBuilder restTemplateBuilder = new RestTemplateBuilder();
+ @PostConstruct
+ public void init() {
+ // Register time module to parse Instant objects.
+ // com.fasterxml.jackson.databind.exc.InvalidDefinitionException:
+ // Java 8 date/time type `java.time.Instant` not supported by default:
+ // add Module "com.fasterxml.jackson.datatype:jackson-datatype-jsr310" to enable handling
+ json.registerModule(new JavaTimeModule());
+ }
+
@Override
public SecurityUser getOrCreateUserByClientPrincipal(HttpServletRequest request, OAuth2AuthenticationToken token, String providerAccessToken, OAuth2Registration registration) {
OAuth2MapperConfig config = registration.getMapperConfig();
diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java
index 303a430e77..3c7eb12ef3 100644
--- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java
+++ b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java
@@ -15,6 +15,7 @@
*/
package org.thingsboard.server.service.security.auth.oauth2;
+import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.security.core.Authentication;
import org.springframework.security.oauth2.client.OAuth2AuthorizedClient;
@@ -42,6 +43,7 @@ import java.net.URLEncoder;
import java.nio.charset.StandardCharsets;
import java.util.UUID;
+@Slf4j
@Component(value = "oauth2AuthenticationSuccessHandler")
public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationSuccessHandler {
@@ -99,6 +101,8 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
clearAuthenticationAttributes(request, response);
getRedirectStrategy().sendRedirect(request, response, baseUrl + "/?accessToken=" + accessToken.getToken() + "&refreshToken=" + refreshToken.getToken());
} catch (Exception e) {
+ log.debug("Error occurred during processing authentication success result. " +
+ "request [{}], response [{}], authentication [{}]", request, response, authentication, e);
clearAuthenticationAttributes(request, response);
String errorPrefix;
if (!StringUtils.isEmpty(callbackUrlScheme)) {
diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
index bf9fb066a0..26b2fe5324 100644
--- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
@@ -137,7 +137,6 @@ public class DefaultDeviceStateService extends TbApplicationEventListener> partitionedDevices = new ConcurrentHashMap<>();
final ConcurrentMap deviceStates = new ConcurrentHashMap<>();
- private final ConcurrentMap deviceLastSavedActivity = new ConcurrentHashMap<>();
final Queue> subscribeQueue = new ConcurrentLinkedQueue<>();
@@ -192,7 +191,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener 0 && lastReportedActivity > lastSavedActivity) {
- final DeviceStateData stateData = getOrFetchDeviceStateData(deviceId);
+ final DeviceStateData stateData = getOrFetchDeviceStateData(deviceId);
+ if (lastReportedActivity > 0 && lastReportedActivity > stateData.getState().getLastActivityTime()) {
updateActivityState(deviceId, stateData, lastReportedActivity);
}
+ cleanDeviceStateIfBelongsExternalPartition(tenantId, deviceId);
}
void updateActivityState(DeviceId deviceId, DeviceStateData stateData, long lastReportedActivity) {
log.trace("updateActivityState - fetched state {} for device {}, lastReportedActivity {}", stateData, deviceId, lastReportedActivity);
if (stateData != null) {
save(deviceId, LAST_ACTIVITY_TIME, lastReportedActivity);
- deviceLastSavedActivity.put(deviceId, lastReportedActivity);
DeviceState state = stateData.getState();
state.setLastActivityTime(lastReportedActivity);
if (!state.isActive()) {
@@ -225,21 +224,23 @@ public class DefaultDeviceStateService extends TbApplicationEventListener {
Set devices = partitionedDevices.remove(partition);
- devices.forEach(deviceId -> {
- deviceStates.remove(deviceId);
- deviceLastSavedActivity.remove(deviceId);
- });
+ devices.forEach(this::cleanUpDeviceStateMap);
});
addedPartitions.forEach(tpi -> partitionedDevices.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet()));
@@ -463,11 +460,12 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIds = new HashSet<>(deviceStates.keySet());
- for (DeviceId deviceId : deviceIds) {
- updateInactivityStateIfExpired(ts, deviceId);
- }
+ partitionedDevices.forEach((tpi, deviceIds) -> {
+ log.debug("Calculating state updates. tpi {} for {} devices", tpi.getFullTopicName(), deviceIds.size());
+ for (DeviceId deviceId : deviceIds) {
+ updateInactivityStateIfExpired(ts, deviceId);
+ }
+ });
}
void updateInactivityStateIfExpired(long ts, DeviceId deviceId) {
@@ -488,8 +486,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIdSet = partitionedDevices.get(tpi);
deviceIdSet.remove(deviceId);
}
+ private void cleanUpDeviceStateMap(DeviceId deviceId) {
+ deviceStates.remove(deviceId);
+ }
+
private ListenableFuture fetchDeviceState(Device device) {
ListenableFuture future;
if (persistToTelemetry) {
diff --git a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
index b1bbb9cdfa..a717769310 100644
--- a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java
@@ -18,6 +18,7 @@ package org.thingsboard.server.service.state;
import org.springframework.context.ApplicationListener;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.common.msg.queue.TbCallback;
@@ -33,13 +34,13 @@ public interface DeviceStateService extends ApplicationListener(attributes))
@@ -269,10 +269,10 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
callback.onSuccess();
}
- private void updateDeviceInactivityTimeout(EntityId entityId, List extends KvEntry> kvEntries) {
+ private void updateDeviceInactivityTimeout(TenantId tenantId, EntityId entityId, List extends KvEntry> kvEntries) {
for (KvEntry kvEntry : kvEntries) {
if (kvEntry.getKey().equals(DefaultDeviceStateService.INACTIVITY_TIMEOUT)) {
- deviceStateService.onDeviceInactivityTimeoutUpdate(new DeviceId(entityId.getId()), kvEntry.getLongValue().orElse(0L));
+ deviceStateService.onDeviceInactivityTimeoutUpdate(tenantId, new DeviceId(entityId.getId()), kvEntry.getLongValue().orElse(0L));
}
}
}
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 0cede00913..696c314379 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -588,6 +588,7 @@ transport:
bind_address: "${MQTT_BIND_ADDRESS:0.0.0.0}"
bind_port: "${MQTT_BIND_PORT:1883}"
timeout: "${MQTT_TIMEOUT:10000}"
+ msg_queue_size_per_device_limit: "${MQTT_MSG_QUEUE_SIZE_PER_DEVICE_LIMIT:100}" # messages await in the queue before device connected state. This limit works on low level before TenantProfileLimits mechanism
netty:
leak_detector_level: "${NETTY_LEAK_DETECTOR_LVL:DISABLED}"
boss_group_thread_count: "${NETTY_BOSS_GROUP_THREADS:1}"
diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
index 1cab815e6a..a7cd4ddf79 100644
--- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
@@ -27,6 +27,7 @@ import org.springframework.mock.web.MockMultipartFile;
import org.springframework.test.web.servlet.request.MockMultipartHttpServletRequestBuilder;
import org.springframework.test.web.servlet.request.MockMvcRequestBuilders;
import org.thingsboard.common.util.JacksonUtil;
+import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@@ -261,7 +262,7 @@ public class AbstractLwM2MIntegrationTest extends AbstractWebsocketTest {
@Before
public void beforeTest() throws Exception {
- executor = Executors.newScheduledThreadPool(10);
+ executor = Executors.newScheduledThreadPool(10, ThingsBoardThreadFactory.forName("test-lwm2m-scheduled"));
loginTenantAdmin();
String[] resources = new String[]{"1.xml", "2.xml", "3.xml", "5.xml", "9.xml"};
diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java
index 6703d79713..d97708fe5b 100644
--- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java
@@ -16,6 +16,8 @@
package org.thingsboard.server.transport.lwm2m;
import com.fasterxml.jackson.core.type.TypeReference;
+import lombok.extern.slf4j.Slf4j;
+import org.junit.After;
import org.junit.Assert;
import org.junit.Test;
import org.thingsboard.server.common.data.Device;
@@ -23,16 +25,18 @@ import org.thingsboard.server.common.data.device.credentials.lwm2m.NoSecClientCr
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus;
-import org.thingsboard.server.common.data.query.EntityKey;
-import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.transport.lwm2m.client.LwM2MTestClient;
import java.util.Arrays;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
+import static org.awaitility.Awaitility.await;
+import static org.hamcrest.Matchers.is;
import static org.thingsboard.rest.client.utils.RestJsonConverter.toTimeseries;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.DOWNLOADED;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.DOWNLOADING;
@@ -43,8 +47,10 @@ import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.UPDA
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.UPDATING;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.VERIFIED;
+@Slf4j
public class NoSecLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest {
+ public static final int TIMEOUT = 30;
private final String OTA_TRANSPORT_CONFIGURATION = "{\n" +
" \"observeAttr\": {\n" +
" \"keyName\": {\n" +
@@ -122,6 +128,15 @@ public class NoSecLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest {
" \"type\": \"LWM2M\"\n" +
"}";
+ LwM2MTestClient client = null;
+
+ @After
+ public void tearDown() {
+ if (client != null) {
+ client.destroy();
+ }
+ }
+
@Test
public void testConnectAndObserveTelemetry() throws Exception {
NoSecClientCredentials clientCredentials = new NoSecClientCredentials();
@@ -196,37 +211,68 @@ public class NoSecLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest {
}
}
+ /**
+ * This is the example how to use the AWAITILITY instead Thread.sleep()
+ * Test will finish as fast as possible, but will await until TIMEOUT if a build machine is busy or slow
+ * Check the detailed log output to learn how Awaitility polling the API and when exactly expected result appears
+ * */
@Test
public void testSoftwareUpdateByObject9() throws Exception {
- LwM2MTestClient client = null;
- try {
- createDeviceProfile(OTA_TRANSPORT_CONFIGURATION);
- NoSecClientCredentials clientCredentials = new NoSecClientCredentials();
- clientCredentials.setEndpoint("OTA_" + ENDPOINT);
- Device device = createDevice(clientCredentials);
-
- device.setSoftwareId(createSoftware().getId());
- device = doPost("/api/device", device, Device.class);
+ //given
+ final List expectedStatuses = Collections.unmodifiableList(Arrays.asList(
+ QUEUED, INITIATED, DOWNLOADING, DOWNLOADING, DOWNLOADING, DOWNLOADED, VERIFIED, UPDATED));
- Thread.sleep(1000);
-
- client = new LwM2MTestClient(executor, "OTA_" + ENDPOINT);
- client.init(SECURITY, COAP_CONFIG);
-
- Thread.sleep(3000);
-
- List ts = toTimeseries(doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + device.getId().getId() + "/values/timeseries?orderBy=ASC&keys=sw_state&startTs=0&endTs=" + System.currentTimeMillis(), new TypeReference<>() {
- }));
-
- List statuses = ts.stream().sorted(Comparator.comparingLong(TsKvEntry::getTs)).map(KvEntry::getValueAsString).map(OtaPackageUpdateStatus::valueOf).collect(Collectors.toList());
+ createDeviceProfile(OTA_TRANSPORT_CONFIGURATION);
+ NoSecClientCredentials clientCredentials = new NoSecClientCredentials();
+ clientCredentials.setEndpoint("OTA_" + ENDPOINT);
+ final Device device = createDevice(clientCredentials);
+ device.setSoftwareId(createSoftware().getId());
+
+ log.warn("Saving by API " + device);
+ final Device savedDevice = doPost("/api/device", device, Device.class);
+ Assert.assertNotNull(savedDevice);
+ log.warn("Device saved by API {}", savedDevice);
+
+ log.warn("AWAIT atMost {} SECONDS on get device by API...", TIMEOUT);
+ await()
+ .atMost(TIMEOUT, TimeUnit.SECONDS)
+ .until(() -> getDeviceFromAPI(device.getId().getId()), is(savedDevice));
+ log.warn("Got device by API.");
+
+ //when
+ log.warn("Init the client...");
+ client = new LwM2MTestClient(executor, "OTA_" + ENDPOINT);
+ client.init(SECURITY, COAP_CONFIG);
+ log.warn("Init done");
+
+ log.warn("AWAIT atMost {} SECONDS on timeseries List by API with list size {}...", TIMEOUT, expectedStatuses.size());
+ await()
+ .atMost(30, TimeUnit.SECONDS)
+ .until(() -> getSwStateTelemetryFromAPI(device.getId().getId())
+ .size(), is(expectedStatuses.size()));
+ log.warn("Got an expected await condition!");
+
+ //then
+ log.warn("Fetching ts for the final asserts");
+ List ts = getSwStateTelemetryFromAPI(device.getId().getId());
+ log.warn("Got an ts {}", ts);
+
+ List statuses = ts.stream().sorted(Comparator.comparingLong(TsKvEntry::getTs)).map(KvEntry::getValueAsString).map(OtaPackageUpdateStatus::valueOf).collect(Collectors.toList());
+ log.warn("Converted ts to statuses {}", statuses);
+
+ Assert.assertEquals(expectedStatuses, statuses);
+ }
- List expectedStatuses = Arrays.asList(QUEUED, INITIATED, DOWNLOADING, DOWNLOADING, DOWNLOADING, DOWNLOADED, VERIFIED, UPDATED);
+ private Device getDeviceFromAPI(UUID deviceId) throws Exception {
+ final Device device = doGet("/api/device/" + deviceId, Device.class);
+ log.warn("Fetched device by API for deviceId {}, device is {}", deviceId, device);
+ return device;
+ }
- Assert.assertEquals(expectedStatuses, statuses);
- } finally {
- if (client != null) {
- client.destroy();
- }
- }
+ private List getSwStateTelemetryFromAPI(UUID deviceId) throws Exception {
+ final List tsKvEntries = toTimeseries(doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + deviceId + "/values/timeseries?orderBy=ASC&keys=sw_state&startTs=0&endTs=" + System.currentTimeMillis(), new TypeReference<>() {
+ }));
+ log.warn("Fetched telemetry by API for deviceId {}, list size {}, tsKvEntries {}", deviceId, tsKvEntries.size(), tsKvEntries);
+ return tsKvEntries;
}
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java
index 8704631bc6..a5aa451bd1 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java
@@ -77,25 +77,34 @@ public class TbKafkaProducerTemplate implements TbQueuePro
@Override
public void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback) {
- createTopicIfNotExist(tpi);
- String key = msg.getKey().toString();
- byte[] data = msg.getData();
- ProducerRecord record;
- Iterable headers = msg.getHeaders().getData().entrySet().stream().map(e -> new RecordHeader(e.getKey(), e.getValue())).collect(Collectors.toList());
- record = new ProducerRecord<>(tpi.getFullTopicName(), null, key, data, headers);
- producer.send(record, (metadata, exception) -> {
- if (exception == null) {
- if (callback != null) {
- callback.onSuccess(new KafkaTbQueueMsgMetadata(metadata));
- }
- } else {
- if (callback != null) {
- callback.onFailure(exception);
+ try {
+ createTopicIfNotExist(tpi);
+ String key = msg.getKey().toString();
+ byte[] data = msg.getData();
+ ProducerRecord record;
+ Iterable headers = msg.getHeaders().getData().entrySet().stream().map(e -> new RecordHeader(e.getKey(), e.getValue())).collect(Collectors.toList());
+ record = new ProducerRecord<>(tpi.getFullTopicName(), null, key, data, headers);
+ producer.send(record, (metadata, exception) -> {
+ if (exception == null) {
+ if (callback != null) {
+ callback.onSuccess(new KafkaTbQueueMsgMetadata(metadata));
+ }
} else {
- log.warn("Producer template failure: {}", exception.getMessage(), exception);
+ if (callback != null) {
+ callback.onFailure(exception);
+ } else {
+ log.warn("Producer template failure: {}", exception.getMessage(), exception);
+ }
}
+ });
+ } catch (Exception e) {
+ if (callback != null) {
+ callback.onFailure(e);
+ } else {
+ log.warn("Producer template failure (send method wrapper): {}", e.getMessage(), e);
}
- });
+ throw e;
+ }
}
private void createTopicIfNotExist(TopicPartitionInfo tpi) {
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java
index c2ba0b853b..16f9443547 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapService.java
@@ -19,6 +19,9 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.elements.util.SslContextUtil;
import org.eclipse.californium.scandium.config.DtlsConnectorConfig;
+import org.eclipse.leshan.core.model.ObjectLoader;
+import org.eclipse.leshan.core.model.ObjectModel;
+import org.eclipse.leshan.core.model.StaticModel;
import org.eclipse.leshan.server.bootstrap.BootstrapSessionManager;
import org.eclipse.leshan.server.californium.bootstrap.LeshanBootstrapServer;
import org.eclipse.leshan.server.californium.bootstrap.LeshanBootstrapServerBuilder;
@@ -26,6 +29,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MBootstrapSecurityStore;
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MInMemoryBootstrapConfigStore;
+import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MInMemoryBootstrapConfigurationAdapter;
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2mDefaultBootstrapSessionManager;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
@@ -38,6 +42,7 @@ import java.security.KeyStoreException;
import java.security.PrivateKey;
import java.security.PublicKey;
import java.security.cert.X509Certificate;
+import java.util.List;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mNetworkConfig.getCoapConfig;
@@ -79,12 +84,14 @@ public class LwM2MTransportBootstrapService {
builder.setCoapConfig(getCoapConfig(bootstrapConfig.getPort(), bootstrapConfig.getSecurePort(), serverConfig));
/* Define model provider (Create Models )*/
+ List models = ObjectLoader.loadDefault();
+ builder.setModel(new StaticModel(models));
/* Create credentials */
this.setServerWithCredentials(builder);
-// /** Set securityStore with new ConfigStore */
-// builder.setConfigStore(lwM2MInMemoryBootstrapConfigStore);
+ /* Set securityStore with new ConfigStore */
+ builder.setConfigStore(new LwM2MInMemoryBootstrapConfigurationAdapter(lwM2MInMemoryBootstrapConfigStore));
/* SecurityStore */
builder.setSecurityStore(lwM2MBootstrapSecurityStore);
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java
index 30ac8e01c3..7dca87458b 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java
@@ -74,15 +74,19 @@ public class LwM2MBootstrapConfig implements Serializable {
configBs.servers.put(0, server0);
/* Security Configuration (object 0) as defined in LWM2M 1.0.x TS. Bootstrap instance = 0 */
this.bootstrapServer.setBootstrapServerIs(true);
- configBs.security.put(0, setServerSecurity(this.bootstrapServer.getHost(), this.bootstrapServer.getPort(), this.bootstrapServer.isBootstrapServerIs(), this.bootstrapServer.getSecurityMode(), this.bootstrapServer.getClientPublicKeyOrId(), this.bootstrapServer.getServerPublicKey(), this.bootstrapServer.getClientSecretKey(), this.bootstrapServer.getServerId()));
+ configBs.security.put(0, setServerSecurity(this.lwm2mServer.getHost(), this.lwm2mServer.getPort(), this.lwm2mServer.getSecurityHost(), this.lwm2mServer.getSecurityPort(), this.bootstrapServer.isBootstrapServerIs(), this.bootstrapServer.getSecurityMode(), this.bootstrapServer.getClientPublicKeyOrId(), this.bootstrapServer.getServerPublicKey(), this.bootstrapServer.getClientSecretKey(), this.bootstrapServer.getServerId()));
/* Security Configuration (object 0) as defined in LWM2M 1.0.x TS. Server instance = 1 */
- configBs.security.put(1, setServerSecurity(this.lwm2mServer.getHost(), this.lwm2mServer.getPort(), this.lwm2mServer.isBootstrapServerIs(), this.lwm2mServer.getSecurityMode(), this.lwm2mServer.getClientPublicKeyOrId(), this.lwm2mServer.getServerPublicKey(), this.lwm2mServer.getClientSecretKey(), this.lwm2mServer.getServerId()));
+ configBs.security.put(1, setServerSecurity(this.lwm2mServer.getHost(), this.lwm2mServer.getPort(), this.lwm2mServer.getSecurityHost(), this.lwm2mServer.getSecurityPort(), this.lwm2mServer.isBootstrapServerIs(), this.lwm2mServer.getSecurityMode(), this.lwm2mServer.getClientPublicKeyOrId(), this.lwm2mServer.getServerPublicKey(), this.lwm2mServer.getClientSecretKey(), this.lwm2mServer.getServerId()));
return configBs;
}
- private BootstrapConfig.ServerSecurity setServerSecurity(String host, Integer port, boolean bootstrapServer, SecurityMode securityMode, String clientPublicKey, String serverPublicKey, String secretKey, int serverId) {
+ private BootstrapConfig.ServerSecurity setServerSecurity(String host, Integer port, String securityHost, Integer securityPort, boolean bootstrapServer, SecurityMode securityMode, String clientPublicKey, String serverPublicKey, String secretKey, int serverId) {
BootstrapConfig.ServerSecurity serverSecurity = new BootstrapConfig.ServerSecurity();
- serverSecurity.uri = "coaps://" + host + ":" + Integer.toString(port);
+ if (securityMode.equals(SecurityMode.NO_SEC)) {
+ serverSecurity.uri = "coap://" + host + ":" + Integer.toString(port);
+ } else {
+ serverSecurity.uri = "coaps://" + securityHost + ":" + Integer.toString(securityPort);
+ }
serverSecurity.bootstrapServer = bootstrapServer;
serverSecurity.securityMode = securityMode;
serverSecurity.publicKeyOrId = setPublicKeyOrId(clientPublicKey, securityMode);
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MInMemoryBootstrapConfigurationAdapter.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MInMemoryBootstrapConfigurationAdapter.java
new file mode 100644
index 0000000000..12325f8c22
--- /dev/null
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MInMemoryBootstrapConfigurationAdapter.java
@@ -0,0 +1,27 @@
+/**
+ * Copyright © 2016-2021 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.lwm2m.bootstrap.secure;
+
+import org.eclipse.leshan.server.bootstrap.BootstrapConfigStore;
+import org.eclipse.leshan.server.bootstrap.BootstrapConfigurationStoreAdapter;
+
+public class LwM2MInMemoryBootstrapConfigurationAdapter extends BootstrapConfigurationStoreAdapter {
+
+ public LwM2MInMemoryBootstrapConfigurationAdapter(BootstrapConfigStore store) {
+ super(store);
+ }
+
+}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MServerBootstrap.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MServerBootstrap.java
index 27d2e8c865..c8a004f52c 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MServerBootstrap.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MServerBootstrap.java
@@ -31,24 +31,31 @@ public class LwM2MServerBootstrap {
String host = "0.0.0.0";
Integer port = 0;
+ String securityHost = "0.0.0.0";
+ Integer securityPort = 0;
SecurityMode securityMode = SecurityMode.NO_SEC;
Integer serverId = 123;
boolean bootstrapServerIs = false;
- public LwM2MServerBootstrap(){};
+ public LwM2MServerBootstrap() {
+ }
+
+ ;
public LwM2MServerBootstrap(LwM2MServerBootstrap bootstrapFromCredential, LwM2MServerBootstrap profileServerBootstrap) {
- this.clientPublicKeyOrId = bootstrapFromCredential.getClientPublicKeyOrId();
- this.clientSecretKey = bootstrapFromCredential.getClientSecretKey();
- this.serverPublicKey = profileServerBootstrap.getServerPublicKey();
- this.clientHoldOffTime = profileServerBootstrap.getClientHoldOffTime();
- this.bootstrapServerAccountTimeout = profileServerBootstrap.getBootstrapServerAccountTimeout();
- this.host = (profileServerBootstrap.getHost().equals("0.0.0.0")) ? "localhost" : profileServerBootstrap.getHost();
- this.port = profileServerBootstrap.getPort();
- this.securityMode = profileServerBootstrap.getSecurityMode();
- this.serverId = profileServerBootstrap.getServerId();
- this.bootstrapServerIs = profileServerBootstrap.bootstrapServerIs;
+ this.clientPublicKeyOrId = bootstrapFromCredential.getClientPublicKeyOrId();
+ this.clientSecretKey = bootstrapFromCredential.getClientSecretKey();
+ this.serverPublicKey = profileServerBootstrap.getServerPublicKey();
+ this.clientHoldOffTime = profileServerBootstrap.getClientHoldOffTime();
+ this.bootstrapServerAccountTimeout = profileServerBootstrap.getBootstrapServerAccountTimeout();
+ this.host = (profileServerBootstrap.getHost().equals("0.0.0.0")) ? "localhost" : profileServerBootstrap.getHost();
+ this.port = profileServerBootstrap.getPort();
+ this.securityHost = (profileServerBootstrap.getSecurityHost().equals("0.0.0.0")) ? "localhost" : profileServerBootstrap.getSecurityHost();
+ this.securityPort = profileServerBootstrap.getSecurityPort();
+ this.securityMode = profileServerBootstrap.getSecurityMode();
+ this.serverId = profileServerBootstrap.getServerId();
+ this.bootstrapServerIs = profileServerBootstrap.bootstrapServerIs;
}
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
index 7537e63462..e29dcdd1dd 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
@@ -93,7 +93,12 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
log.debug("Fetched clients from store: {}", fetchedClients);
fetchedClients.forEach(client -> {
lwM2mClientsByEndpoint.put(client.getEndpoint(), client);
- updateFetchedClient(nodeId, client);
+ try {
+ client.lock();
+ updateFetchedClient(nodeId, client);
+ } finally {
+ client.unlock();
+ }
});
}
@@ -161,7 +166,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
this.lwM2mClientsByRegistrationId.put(registration.getId(), client);
client.setState(LwM2MClientState.REGISTERED);
onUplink(client);
- if(!compareAndSetSleepFlag(client, false)){
+ if (!compareAndSetSleepFlag(client, false)) {
clientStore.put(client);
}
} finally {
@@ -311,7 +316,11 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
public void update(LwM2mClient client) {
client.lock();
try {
- clientStore.put(client);
+ if (client.getState().equals(LwM2MClientState.REGISTERED)) {
+ clientStore.put(client);
+ } else {
+ log.error("[{}] Client is in invalid state: {}!", client.getEndpoint(), client.getState());
+ }
} finally {
client.unlock();
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java
index 7261e5bd40..039629516f 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java
@@ -106,6 +106,7 @@ public abstract class LwM2MClientOtaInfo {
public abstract OtaPackageType getType();
+ @JsonIgnore
public String getTargetPackageId() {
return getPackageId(targetName, targetVersion);
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java
index d735eed26e..2a111a9f81 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java
@@ -15,11 +15,13 @@
*/
package org.thingsboard.server.transport.lwm2m.server.store;
+import lombok.extern.slf4j.Slf4j;
import org.nustaq.serialization.FSTConfiguration;
import org.springframework.data.redis.connection.RedisClusterConnection;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.Cursor;
import org.springframework.data.redis.core.ScanOptions;
+import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import java.util.ArrayList;
@@ -27,6 +29,7 @@ import java.util.HashSet;
import java.util.List;
import java.util.Set;
+@Slf4j
public class TbRedisLwM2MClientStore implements TbLwM2MClientStore {
private static final String CLIENT_EP = "CLIENT#EP#";
@@ -76,9 +79,13 @@ public class TbRedisLwM2MClientStore implements TbLwM2MClientStore {
@Override
public void put(LwM2mClient client) {
- byte[] clientSerialized = serializer.asByteArray(client);
- try (var connection = connectionFactory.getConnection()) {
- connection.getSet(getKey(client.getEndpoint()), clientSerialized);
+ if (client.getState().equals(LwM2MClientState.UNREGISTERED)) {
+ log.error("[{}] Client is in invalid state: {}!", client.getEndpoint(), client.getState(), new Exception());
+ } else {
+ byte[] clientSerialized = serializer.asByteArray(client);
+ try (var connection = connectionFactory.getConnection()) {
+ connection.getSet(getKey(client.getEndpoint()), clientSerialized);
+ }
}
}
diff --git a/common/transport/mqtt/pom.xml b/common/transport/mqtt/pom.xml
index 676593804e..b6ceb13951 100644
--- a/common/transport/mqtt/pom.xml
+++ b/common/transport/mqtt/pom.xml
@@ -88,6 +88,11 @@
junit
test
+
+ org.awaitility
+ awaitility
+ test
+
org.mockito
mockito-core
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java
index f160936c7c..11b46696da 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java
@@ -23,10 +23,15 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
+import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor;
import org.thingsboard.server.transport.mqtt.adaptors.ProtoMqttAdaptor;
+import javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
+import java.util.concurrent.ExecutorService;
+
/**
* Created by ashvayka on 04.10.18.
*/
@@ -59,4 +64,8 @@ public class MqttTransportContext extends TransportContext {
@Setter
private SslHandler sslHandler;
+ @Getter
+ @Value("${transport.mqtt.msg_queue_size_per_device_limit:100}")
+ private int messageQueueSizePerDeviceLimit;
+
}
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
index 630aec946d..3e3f2e74d9 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
@@ -123,9 +123,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final SslHandler sslHandler;
private final ConcurrentMap mqttQoSMap;
- private final DeviceSessionCtx deviceSessionCtx;
- private volatile InetSocketAddress address;
- private volatile GatewaySessionHandler gatewaySessionHandler;
+ final DeviceSessionCtx deviceSessionCtx;
+ volatile InetSocketAddress address;
+ volatile GatewaySessionHandler gatewaySessionHandler;
private final ConcurrentHashMap otaPackSessions;
private final ConcurrentHashMap chunkSizes;
@@ -164,8 +164,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
- private void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) {
- address = (InetSocketAddress) ctx.channel().remoteAddress();
+ void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) {
+ address = getAddress(ctx);
if (msg.fixedHeader() == null) {
log.info("[{}:{}] Invalid message received", address.getHostName(), address.getPort());
processDisconnect(ctx);
@@ -177,10 +177,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} else if (deviceSessionCtx.isProvisionOnly()) {
processProvisionSessionMsg(ctx, msg);
} else {
- processRegularSessionMsg(ctx, msg);
+ enqueueRegularSessionMsg(ctx, msg);
}
}
+ InetSocketAddress getAddress(ChannelHandlerContext ctx) {
+ return (InetSocketAddress) ctx.channel().remoteAddress();
+ }
+
private void processProvisionSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) {
switch (msg.fixedHeader().messageType()) {
case PUBLISH:
@@ -223,7 +227,42 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
- private void processRegularSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) {
+ void enqueueRegularSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) {
+ final int queueSize = deviceSessionCtx.getMsgQueueSize().incrementAndGet();
+ if (queueSize > context.getMessageQueueSizePerDeviceLimit()) {
+ log.warn("Closing current session because msq queue size for device {} exceed limit {} with msgQueueSize counter {} and actual queue size {}",
+ deviceSessionCtx.getDeviceId(), context.getMessageQueueSizePerDeviceLimit(), queueSize, deviceSessionCtx.getMsgQueue().size());
+ ctx.close();
+ return;
+ }
+
+ deviceSessionCtx.getMsgQueue().add(msg);
+ processMsgQueue(ctx); //Under the normal conditions the msg queue will contain 0 messages. Many messages will be processed on device connect event in separate thread pool
+ }
+
+ void processMsgQueue(ChannelHandlerContext ctx) {
+ if (!deviceSessionCtx.isConnected()) {
+ log.trace("[{}][{}] Postpone processing msg due to device is not connected. Msg queue size is {}", sessionId, deviceSessionCtx.getDeviceId(), deviceSessionCtx.getMsgQueue().size());
+ return;
+ }
+ while (!deviceSessionCtx.getMsgQueue().isEmpty()) {
+ if (deviceSessionCtx.getMsgQueueProcessorLock().tryLock()) {
+ try {
+ MqttMessage msg;
+ while ((msg = deviceSessionCtx.getMsgQueue().poll()) != null) {
+ deviceSessionCtx.getMsgQueueSize().decrementAndGet();
+ processRegularSessionMsg(ctx, msg);
+ }
+ } finally {
+ deviceSessionCtx.getMsgQueueProcessorLock().unlock();
+ }
+ } else {
+ return;
+ }
+ }
+ }
+
+ void processRegularSessionMsg(ChannelHandlerContext ctx, MqttMessage msg) {
switch (msg.fixedHeader().messageType()) {
case PUBLISH:
processPublish(ctx, (MqttPublishMessage) msg);
@@ -304,6 +343,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
} catch (RuntimeException | AdaptorException e) {
log.warn("[{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId, e);
+ ctx.close();
}
}
@@ -588,7 +628,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
return new MqttMessage(mqttFixedHeader, mqttMessageIdVariableHeader);
}
- private void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) {
+ void processConnect(ChannelHandlerContext ctx, MqttConnectMessage msg) {
log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier());
String userName = msg.payload().userName();
String clientId = msg.payload().clientIdentifier();
@@ -674,7 +714,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
return null;
}
- private void processDisconnect(ChannelHandlerContext ctx) {
+ void processDisconnect(ChannelHandlerContext ctx) {
ctx.close();
log.info("[{}] Client disconnected!", sessionId);
doDisconnect();
@@ -761,6 +801,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
deviceSessionCtx.setDisconnected();
}
+
+ if (!deviceSessionCtx.getMsgQueue().isEmpty()) {
+ log.warn("doDisconnect for device {} but unprocessed messages {} left in the msg queue", deviceSessionCtx.getDeviceId(), deviceSessionCtx.getMsgQueue().size());
+ deviceSessionCtx.getMsgQueue().clear();
+ }
}
@@ -778,7 +823,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
SessionMetaData sessionMetaData = transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this);
checkGatewaySession(sessionMetaData);
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, connectMessage));
+ deviceSessionCtx.setConnected(true);
log.info("[{}] Client connected!", sessionId);
+ transportService.getCallbackExecutor().execute(() -> processMsgQueue(ctx)); //this callback will execute in Producer worker thread and hard or blocking work have to be submitted to the separate thread.
}
@Override
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
index 3804d96cf6..b92e91981f 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
@@ -18,6 +18,7 @@ package org.thingsboard.server.transport.mqtt.session;
import com.google.protobuf.Descriptors;
import com.google.protobuf.DynamicMessage;
import io.netty.channel.ChannelHandlerContext;
+import io.netty.handler.codec.mqtt.MqttMessage;
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
@@ -35,8 +36,11 @@ import org.thingsboard.server.transport.mqtt.util.MqttTopicFilter;
import org.thingsboard.server.transport.mqtt.util.MqttTopicFilterFactory;
import java.util.UUID;
+import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
/**
* @author Andrew Shvayka
@@ -45,13 +49,23 @@ import java.util.concurrent.atomic.AtomicInteger;
public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
@Getter
+ @Setter
private ChannelHandlerContext channel;
@Getter
- private MqttTransportContext context;
+ private final MqttTransportContext context;
private final AtomicInteger msgIdSeq = new AtomicInteger(0);
+ @Getter
+ private final ConcurrentLinkedQueue msgQueue = new ConcurrentLinkedQueue<>();
+
+ @Getter
+ private final Lock msgQueueProcessorLock = new ReentrantLock();
+
+ @Getter
+ private final AtomicInteger msgQueueSize = new AtomicInteger(0);
+
@Getter
@Setter
private boolean provisionOnly = false;
@@ -73,10 +87,6 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
this.context = context;
}
- public void setChannel(ChannelHandlerContext channel) {
- this.channel = channel;
- }
-
public int nextMsgId() {
return msgIdSeq.incrementAndGet();
}
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
index 3fed5e51ea..21086d75d8 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
@@ -60,6 +60,7 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple
.setDeviceProfileIdLSB(deviceInfo.getDeviceProfileId().getId().getLeastSignificantBits())
.build());
setDeviceInfo(deviceInfo);
+ setConnected(true);
setDeviceProfile(deviceProfile);
this.transportService = transportService;
}
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
index 193a7d24f5..0a547eb11a 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
@@ -34,6 +34,7 @@ import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.CollectionUtils;
+import org.springframework.util.ConcurrentReferenceHashMap;
import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.transport.TransportService;
@@ -66,6 +67,8 @@ import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
+import static org.springframework.util.ConcurrentReferenceHashMap.ReferenceType;
+
/**
* Created by ashvayka on 19.01.17.
*/
@@ -82,7 +85,7 @@ public class GatewaySessionHandler {
private final UUID sessionId;
private final ConcurrentMap deviceCreationLockMap;
private final ConcurrentMap devices;
- private final ConcurrentMap> deviceFutures;
+ private final ConcurrentMap> deviceFutures;
private final ConcurrentMap mqttQoSMap;
private final ChannelHandlerContext channel;
private final DeviceSessionCtx deviceSessionCtx;
@@ -95,11 +98,15 @@ public class GatewaySessionHandler {
this.sessionId = sessionId;
this.devices = new ConcurrentHashMap<>();
this.deviceFutures = new ConcurrentHashMap<>();
- this.deviceCreationLockMap = new ConcurrentHashMap<>();
+ this.deviceCreationLockMap = createWeakMap();
this.mqttQoSMap = deviceSessionCtx.getMqttQoSMap();
this.channel = deviceSessionCtx.getChannel();
}
+ ConcurrentReferenceHashMap createWeakMap() {
+ return new ConcurrentReferenceHashMap<>(16, ReferenceType.WEAK);
+ }
+
public void onDeviceConnect(MqttPublishMessage mqttMsg) throws AdaptorException {
if (isJsonPayloadType()) {
onDeviceConnectJson(mqttMsg);
@@ -228,21 +235,22 @@ public class GatewaySessionHandler {
if (result == null) {
return getDeviceCreationFuture(deviceName, deviceType);
} else {
- return toCompletedFuture(result);
+ return Futures.immediateFuture(result);
}
} finally {
deviceCreationLock.unlock();
}
} else {
- return toCompletedFuture(result);
+ return Futures.immediateFuture(result);
}
}
private ListenableFuture getDeviceCreationFuture(String deviceName, String deviceType) {
- SettableFuture future = deviceFutures.get(deviceName);
- if (future == null) {
- final SettableFuture futureToSet = SettableFuture.create();
- deviceFutures.put(deviceName, futureToSet);
+ final SettableFuture futureToSet = SettableFuture.create();
+ ListenableFuture future = deviceFutures.putIfAbsent(deviceName, futureToSet);
+ if (future != null) {
+ return future;
+ }
try {
transportService.process(GetOrCreateDeviceFromGatewayRequestMsg.newBuilder()
.setDeviceName(deviceName)
@@ -282,15 +290,6 @@ public class GatewaySessionHandler {
deviceFutures.remove(deviceName);
throw e;
}
- } else {
- return future;
- }
- }
-
- private ListenableFuture toCompletedFuture(GatewayDeviceSessionCtx result) {
- SettableFuture future = SettableFuture.create();
- future.set(result);
- return future;
}
private int getMsgId(MqttPublishMessage mqttMsg) {
@@ -353,6 +352,7 @@ public class GatewaySessionHandler {
processPostTelemetryMsg(deviceCtx, postTelemetryMsg, deviceName, msgId);
} catch (Throwable e) {
log.warn("[{}][{}] Failed to convert telemetry: {}", gateway.getDeviceId(), deviceName, deviceEntry.getValue(), e);
+ channel.close();
}
}
@@ -384,6 +384,7 @@ public class GatewaySessionHandler {
processPostTelemetryMsg(deviceCtx, postTelemetryMsg, deviceName, msgId);
} catch (Throwable e) {
log.warn("[{}][{}] Failed to convert telemetry: {}", gateway.getDeviceId(), deviceName, msg, e);
+ channel.close();
}
}
diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java
new file mode 100644
index 0000000000..9b6367da30
--- /dev/null
+++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java
@@ -0,0 +1,221 @@
+/**
+ * Copyright © 2016-2021 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 io.netty.buffer.ByteBuf;
+import io.netty.buffer.EmptyByteBuf;
+import io.netty.buffer.PooledByteBufAllocator;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.handler.codec.mqtt.MqttConnectMessage;
+import io.netty.handler.codec.mqtt.MqttConnectPayload;
+import io.netty.handler.codec.mqtt.MqttConnectVariableHeader;
+import io.netty.handler.codec.mqtt.MqttFixedHeader;
+import io.netty.handler.codec.mqtt.MqttMessage;
+import io.netty.handler.codec.mqtt.MqttMessageType;
+import io.netty.handler.codec.mqtt.MqttPublishMessage;
+import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
+import io.netty.handler.codec.mqtt.MqttQoS;
+import io.netty.handler.ssl.SslHandler;
+import lombok.extern.slf4j.Slf4j;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.Mock;
+import org.mockito.junit.MockitoJUnitRunner;
+import org.thingsboard.common.util.ThingsBoardThreadFactory;
+
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.contains;
+import static org.hamcrest.Matchers.empty;
+import static org.hamcrest.Matchers.greaterThan;
+import static org.hamcrest.Matchers.is;
+import static org.junit.Assert.fail;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.BDDMockito.willDoNothing;
+import static org.mockito.BDDMockito.willReturn;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+@Slf4j
+@RunWith(MockitoJUnitRunner.class)
+public class MqttTransportHandlerTest {
+
+ public static final int MSG_QUEUE_LIMIT = 10;
+ public static final InetSocketAddress IP_ADDR = new InetSocketAddress("127.0.0.1", 9876);
+ public static final int TIMEOUT = 30;
+
+ @Mock
+ MqttTransportContext context;
+ @Mock
+ SslHandler sslHandler;
+ @Mock
+ ChannelHandlerContext ctx;
+
+ AtomicInteger packedId = new AtomicInteger();
+ ExecutorService executor;
+ MqttTransportHandler handler;
+
+ @Before
+ public void setUp() throws Exception {
+ willReturn(MSG_QUEUE_LIMIT).given(context).getMessageQueueSizePerDeviceLimit();
+
+ handler = spy(new MqttTransportHandler(context, sslHandler));
+ willReturn(IP_ADDR).given(handler).getAddress(any());
+ }
+
+ @After
+ public void tearDown() {
+ if (executor != null) {
+ executor.shutdownNow();
+ }
+ }
+
+ MqttConnectMessage getMqttConnectMessage() {
+ MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.CONNECT, true, MqttQoS.AT_LEAST_ONCE, false, 123);
+ MqttConnectVariableHeader variableHeader = new MqttConnectVariableHeader("device", packedId.incrementAndGet(), true, true, true, 1, true, false, 60);
+ MqttConnectPayload payload = new MqttConnectPayload("clientId", "topic", "message".getBytes(StandardCharsets.UTF_8), "username", "password".getBytes(StandardCharsets.UTF_8));
+ return new MqttConnectMessage(mqttFixedHeader, variableHeader, payload);
+ }
+
+ MqttPublishMessage getMqttPublishMessage() {
+ MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, true, MqttQoS.AT_LEAST_ONCE, false, 123);
+ MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader("v1/gateway/telemetry", packedId.incrementAndGet());
+ ByteBuf payload = new EmptyByteBuf(new PooledByteBufAllocator());
+ return new MqttPublishMessage(mqttFixedHeader, variableHeader, payload);
+ }
+
+ @Test
+ public void givenMessageWithoutFixedHeader_whenProcessMqttMsg_thenProcessDisconnect() {
+ MqttFixedHeader mqttFixedHeader = null;
+ MqttMessage msg = new MqttMessage(mqttFixedHeader);
+ willDoNothing().given(handler).processDisconnect(ctx);
+
+ handler.processMqttMsg(ctx, msg);
+
+ assertThat(handler.address, is(IP_ADDR));
+ verify(handler, times(1)).processDisconnect(ctx);
+ }
+
+ @Test
+ public void givenMqttConnectMessage_whenProcessMqttMsg_thenProcessConnect() {
+ MqttConnectMessage msg = getMqttConnectMessage();
+ willDoNothing().given(handler).processConnect(ctx, msg);
+
+ handler.processMqttMsg(ctx, msg);
+
+ assertThat(handler.address, is(IP_ADDR));
+ assertThat(handler.deviceSessionCtx.getChannel(), is(ctx));
+ verify(handler, never()).processDisconnect(any());
+ verify(handler, times(1)).processConnect(ctx, msg);
+ }
+
+ @Test
+ public void givenQueueLimit_whenEnqueueRegularSessionMsgOverLimit_thenOK() {
+ List messages = Stream.generate(this::getMqttPublishMessage).limit(MSG_QUEUE_LIMIT).collect(Collectors.toList());
+ messages.forEach(msg -> handler.enqueueRegularSessionMsg(ctx, msg));
+ assertThat(handler.deviceSessionCtx.getMsgQueueSize().get(), is(MSG_QUEUE_LIMIT));
+ assertThat(handler.deviceSessionCtx.getMsgQueue(), contains(messages.toArray()));
+ }
+
+ @Test
+ public void givenQueueLimit_whenEnqueueRegularSessionMsgOverLimit_thenCtxClose() {
+ final int limit = MSG_QUEUE_LIMIT + 1;
+ willDoNothing().given(handler).processMsgQueue(ctx);
+ List messages = Stream.generate(this::getMqttPublishMessage).limit(limit).collect(Collectors.toList());
+
+ messages.forEach((msg) -> handler.enqueueRegularSessionMsg(ctx, msg));
+
+ assertThat(handler.deviceSessionCtx.getMsgQueueSize().get(), is(limit));
+ verify(handler, times(limit)).enqueueRegularSessionMsg(any(), any());
+ verify(handler, times(MSG_QUEUE_LIMIT)).processMsgQueue(any());
+ verify(ctx, times(1)).close();
+ }
+
+ @Test
+ public void givenMqttConnectMessageAndPublishImmediately_whenProcessMqttMsg_thenEnqueueRegularSessionMsg() {
+ givenMqttConnectMessage_whenProcessMqttMsg_thenProcessConnect();
+
+ List messages = Stream.generate(this::getMqttPublishMessage).limit(MSG_QUEUE_LIMIT).collect(Collectors.toList());
+
+ messages.forEach((msg) -> handler.processMqttMsg(ctx, msg));
+
+ assertThat(handler.address, is(IP_ADDR));
+ assertThat(handler.deviceSessionCtx.getChannel(), is(ctx));
+ assertThat(handler.deviceSessionCtx.isConnected(), is(false));
+ assertThat(handler.deviceSessionCtx.getMsgQueueSize().get(), is(MSG_QUEUE_LIMIT));
+ assertThat(handler.deviceSessionCtx.getMsgQueue(), contains(messages.toArray()));
+ verify(handler, never()).processDisconnect(any());
+ verify(handler, times(1)).processConnect(any(), any());
+ verify(handler, times(MSG_QUEUE_LIMIT)).enqueueRegularSessionMsg(any(), any());
+ verify(handler, never()).processRegularSessionMsg(any(), any());
+ messages.forEach((msg) -> verify(handler, times(1)).enqueueRegularSessionMsg(ctx, msg));
+ }
+
+ @Test
+ public void givenMessageQueue_whenProcessMqttMsgConcurrently_thenEnqueueRegularSessionMsg() throws InterruptedException {
+ //given
+ assertThat(handler.deviceSessionCtx.isConnected(), is(false));
+ assertThat(MSG_QUEUE_LIMIT, greaterThan(2));
+ List messages = Stream.generate(this::getMqttPublishMessage).limit(MSG_QUEUE_LIMIT).collect(Collectors.toList());
+ messages.forEach((msg) -> handler.enqueueRegularSessionMsg(ctx, msg));
+ willDoNothing().given(handler).processRegularSessionMsg(any(), any());
+ executor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getClass().getName()));
+
+ CountDownLatch readyLatch = new CountDownLatch(MSG_QUEUE_LIMIT);
+ CountDownLatch startLatch = new CountDownLatch(1);
+ CountDownLatch finishLatch = new CountDownLatch(MSG_QUEUE_LIMIT);
+
+ Stream.iterate(0, i -> i + 1).limit(MSG_QUEUE_LIMIT).forEach(x ->
+ executor.submit(() -> {
+ try {
+ readyLatch.countDown();
+ assertThat(startLatch.await(TIMEOUT, TimeUnit.SECONDS), is(true));
+ handler.processMsgQueue(ctx);
+ finishLatch.countDown();
+ } catch (Exception e) {
+ log.error("Failed to run processMsgQueue", e);
+ fail("Failed to run processMsgQueue");
+ }
+ }));
+
+ //when
+ assertThat(readyLatch.await(TIMEOUT, TimeUnit.SECONDS), is(true));
+ handler.deviceSessionCtx.setConnected(true);
+ startLatch.countDown();
+ assertThat(finishLatch.await(TIMEOUT, TimeUnit.SECONDS), is(true));
+
+ //then
+ assertThat(handler.deviceSessionCtx.getMsgQueueSize().get(), is(0));
+ assertThat(handler.deviceSessionCtx.getMsgQueue(), empty());
+ verify(handler, times(MSG_QUEUE_LIMIT)).processRegularSessionMsg(any(), any());
+ messages.forEach((msg) -> verify(handler, times(1)).processRegularSessionMsg(ctx, msg));
+ }
+
+}
\ No newline at end of file
diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandlerTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandlerTest.java
new file mode 100644
index 0000000000..ed45e9e44e
--- /dev/null
+++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandlerTest.java
@@ -0,0 +1,61 @@
+/**
+ * Copyright © 2016-2021 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.session;
+
+import org.junit.Test;
+
+import java.util.WeakHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.BDDMockito.willCallRealMethod;
+import static org.mockito.Mockito.mock;
+
+public class GatewaySessionHandlerTest {
+
+ @Test
+ public void givenWeakHashMap_WhenGC_thenMapIsEmpty() {
+ WeakHashMap map = new WeakHashMap<>();
+
+ String deviceName = new String("device"); //constants are static and doesn't affected by GC, so use new instead
+ map.put(deviceName, new ReentrantLock());
+ assertTrue(map.containsKey(deviceName));
+
+ deviceName = null;
+ System.gc();
+
+ await().atMost(10, TimeUnit.SECONDS).until(() -> !map.containsKey("device"));
+ }
+
+ @Test
+ public void givenConcurrentReferenceHashMap_WhenGC_thenMapIsEmpty() {
+ GatewaySessionHandler gsh = mock(GatewaySessionHandler.class);
+ willCallRealMethod().given(gsh).createWeakMap();
+
+ ConcurrentMap map = gsh.createWeakMap();
+ map.put("device", new ReentrantLock());
+ assertTrue(map.containsKey("device"));
+
+ System.gc();
+
+ await().atMost(10, TimeUnit.SECONDS).until(() -> !map.containsKey("device"));
+ }
+
+}
\ No newline at end of file
diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
index 0efde4d893..bc234340cb 100644
--- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
+++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
@@ -188,6 +188,7 @@ public class SnmpTransportContext extends TransportContext {
deviceSessionContext.setSessionInfo(sessionInfo);
deviceSessionContext.setDeviceInfo(msg.getDeviceInfo());
+ deviceSessionContext.setConnected(true);
} else {
log.warn("[{}] Failed to process device auth", deviceSessionContext.getDeviceId());
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
index 5f2fa4f197..237954c553 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
@@ -56,6 +56,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceLwM2MC
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceTokenRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509CertRequestMsg;
+import java.util.concurrent.ExecutorService;
+
/**
* Created by ashvayka on 04.10.18.
*/
@@ -131,4 +133,6 @@ public interface TransportService {
void log(SessionInfoProto sessionInfo, String msg);
void notifyAboutUplink(SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg build, TransportServiceCallback empty);
+
+ ExecutorService getCallbackExecutor();
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
index e3390eb9bb..c705d8364a 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
@@ -1141,4 +1141,9 @@ public class DefaultTransportService implements TransportService {
callback.onError(e);
}
}
+
+ @Override
+ public ExecutorService getCallbackExecutor() {
+ return transportCallbackExecutor;
+ }
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java
index 7d5bc281ec..1d2b7382d0 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java
@@ -46,6 +46,7 @@ public abstract class DeviceAwareSessionContext implements SessionContext {
@Setter
private volatile TransportProtos.SessionInfoProto sessionInfo;
+ @Setter
private volatile boolean connected;
public DeviceId getDeviceId() {
@@ -54,7 +55,6 @@ public abstract class DeviceAwareSessionContext implements SessionContext {
public void setDeviceInfo(TransportDeviceInfo deviceInfo) {
this.deviceInfo = deviceInfo;
- this.connected = true;
this.deviceId = deviceInfo.getDeviceId();
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java
index 378e7d03bb..13a67c0161 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java
@@ -52,7 +52,7 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
public AdminSettings saveAdminSettings(TenantId tenantId, AdminSettings adminSettings) {
log.trace("Executing saveAdminSettings [{}]", adminSettings);
adminSettingsValidator.validate(adminSettings, data -> tenantId);
- if (adminSettings.getKey().equals("mail") && "".equals(adminSettings.getJsonValue().get("password").asText())) {
+ if(adminSettings.getKey().equals("mail") && !adminSettings.getJsonValue().has("password")) {
AdminSettings mailSettings = findAdminSettingsByKey(tenantId, "mail");
if (mailSettings != null) {
((ObjectNode) adminSettings.getJsonValue()).put("password", mailSettings.getJsonValue().get("password").asText());
@@ -61,7 +61,7 @@ public class AdminSettingsServiceImpl implements AdminSettingsService {
return adminSettingsDao.save(tenantId, adminSettings);
}
-
+
private DataValidator adminSettingsValidator =
new DataValidator() {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
index bca52d3ad4..9c2dce109b 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
@@ -100,7 +100,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
}
@Override
- public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) {
+ public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key) {
return Futures.immediateFuture(null);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
index 31f3407fbf..435c9c5b70 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
@@ -124,7 +124,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
}
@Override
- public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) {
+ public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key) {
return Futures.immediateFuture(0);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
index fb15af723a..3170f6c256 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
@@ -170,7 +170,7 @@ public class BaseTimeseriesService implements TimeseriesService {
if (entityId.getEntityType().equals(EntityType.ENTITY_VIEW)) {
throw new IncorrectParameterException("Telemetry data can't be stored for entity view. Read only");
}
- futures.add(timeseriesDao.savePartition(tenantId, entityId, tsKvEntry.getTs(), tsKvEntry.getKey(), ttl));
+ futures.add(timeseriesDao.savePartition(tenantId, entityId, tsKvEntry.getTs(), tsKvEntry.getKey()));
futures.add(Futures.transform(timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry), v -> 0, MoreExecutors.directExecutor()));
futures.add(timeseriesDao.save(tenantId, entityId, tsKvEntry, ttl));
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
index ce653e2e6e..736b234eb9 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
@@ -181,11 +181,14 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
}
@Override
- public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) {
+ public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key) {
if (isFixedPartitioning()) {
return Futures.immediateFuture(null);
}
- ttl = computeTtl(ttl);
+ // DO NOT apply custom TTL to partition, otherwise, short TTL will remove partition too early
+ // partitions must remain in the DB forever or be removed only by systemTtl
+ // removal of empty partition is too expensive (we need to scan all data keys for these partitions with ALLOW FILTERING)
+ long ttl = computeTtl(0);
long partition = toPartitionTs(tsKvEntryTs);
if (cassandraTsPartitionsCache == null) {
return doSavePartition(tenantId, entityId, key, ttl, partition);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
index e9af5f0b75..5700410fbb 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
@@ -33,7 +33,7 @@ public interface TimeseriesDao {
ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl);
- ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl);
+ ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key);
ListenableFuture remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query);
diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java
index d3c6c97367..e8eae60131 100644
--- a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java
+++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java
@@ -100,10 +100,10 @@ public class CassandraPartitionsCacheTest {
long tsKvEntryTs = System.currentTimeMillis();
for (int i = 0; i < 50000; i++) {
- cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0);
+ cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i);
}
for (int i = 0; i < 60000; i++) {
- cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0);
+ cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i);
}
verify(cassandraBaseTimeseriesDao, times(60000)).executeAsyncWrite(any(TenantId.class), any(Statement.class));
}
diff --git a/pom.xml b/pom.xml
index d0c9864ef3..59c160c6a9 100755
--- a/pom.xml
+++ b/pom.xml
@@ -42,13 +42,14 @@
2.3.12.RELEASE
5.2.16.RELEASE
5.2.11.RELEASE
- 5.4.1
+ 5.4.4
2.4.3
3.3.0
0.7.0
2.2.0
4.12
5.7.1
+ 4.1.0
2.2
1.7.7
1.2.3
@@ -1382,6 +1383,12 @@
io.grpc
grpc-netty
${grpc.version}
+
+
+ io.netty
+ *
+
+
io.grpc
@@ -1437,6 +1444,12 @@
${junit.version}
test
+
+ org.awaitility
+ awaitility
+ ${awaitility.version}
+ test
+
org.hamcrest
hamcrest
@@ -1580,6 +1593,12 @@
com.microsoft.azure
azure-servicebus
${azure-servicebus.version}
+
+
+ io.netty
+ *
+
+
org.passay
diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java
index b0730a5841..2bf9122934 100644
--- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java
+++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java
@@ -188,7 +188,7 @@ class AlarmState {
setAlarmConditionMetadata(ruleState, metaData);
TbMsg newMsg = ctx.newMsg(lastMsgQueueName != null ? lastMsgQueueName : ServiceQueue.MAIN, "ALARM",
originator, msg != null ? msg.getCustomerId() : null, metaData, data);
- ctx.tellNext(newMsg, relationType);
+ ctx.enqueueForTellNext(newMsg, relationType);
}
protected void setAlarmConditionMetadata(AlarmRuleState ruleState, TbMsgMetaData metaData) {
diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
index 3d66c569d8..1018d70a2a 100644
--- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
+++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
@@ -120,7 +120,7 @@ public class TbSendRPCRequestNode implements TbNode {
ctx.enqueueForTellNext(next, TbRelationTypes.SUCCESS);
} else {
TbMsg next = ctx.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getCustomerId(), msg.getMetaData(), wrap("error", ruleEngineDeviceRpcResponse.getError().get().name()));
- ctx.tellFailure(next, new RuntimeException(ruleEngineDeviceRpcResponse.getError().get().name()));
+ ctx.enqueueForTellFailure(next, ruleEngineDeviceRpcResponse.getError().get().name());
}
});
ctx.ack(msg);
diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNodeTest.java
index c3bbc8bcf6..b2c79102c4 100644
--- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNodeTest.java
+++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNodeTest.java
@@ -196,7 +196,7 @@ public class TbDeviceProfileNodeTest {
TbMsgDataType.JSON, mapper.writeValueAsString(data), null, null);
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
TbMsg theMsg2 = TbMsg.newMsg("ALARM", deviceId, new TbMsgMetaData(), "2");
@@ -207,7 +207,7 @@ public class TbDeviceProfileNodeTest {
TbMsgDataType.JSON, mapper.writeValueAsString(data), null, null);
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg2, "Alarm Updated");
+ verify(ctx).enqueueForTellNext(theMsg2, "Alarm Updated");
}
@@ -289,7 +289,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -376,7 +376,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -445,7 +445,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -554,7 +554,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -678,7 +678,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -781,7 +781,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -897,7 +897,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -999,7 +999,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg2);
verify(ctx).tellSuccess(msg2);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -1082,7 +1082,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -1163,7 +1163,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -1237,7 +1237,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -1321,7 +1321,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
@@ -1407,7 +1407,7 @@ public class TbDeviceProfileNodeTest {
node.onMsg(ctx, msg);
verify(ctx).tellSuccess(msg);
- verify(ctx).tellNext(theMsg, "Alarm Created");
+ verify(ctx).enqueueForTellNext(theMsg, "Alarm Created");
verify(ctx, Mockito.never()).tellFailure(Mockito.any(), Mockito.any());
}
diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
index c678d5afed..000acab007 100644
--- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
+++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
@@ -89,6 +89,7 @@ transport:
bind_address: "${MQTT_BIND_ADDRESS:0.0.0.0}"
bind_port: "${MQTT_BIND_PORT:1883}"
timeout: "${MQTT_TIMEOUT:10000}"
+ msg_queue_size_per_device_limit: "${MQTT_MSG_QUEUE_SIZE_PER_DEVICE_LIMIT:100}" # messages await in the queue before device connected state. This limit works on low level before TenantProfileLimits mechanism
netty:
leak_detector_level: "${NETTY_LEAK_DETECTOR_LVL:DISABLED}"
boss_group_thread_count: "${NETTY_BOSS_GROUP_THREADS:1}"
diff --git a/ui-ngx/src/app/core/api/widget-api.models.ts b/ui-ngx/src/app/core/api/widget-api.models.ts
index d473296a4c..9f99e7aed5 100644
--- a/ui-ngx/src/app/core/api/widget-api.models.ts
+++ b/ui-ngx/src/app/core/api/widget-api.models.ts
@@ -84,6 +84,8 @@ export interface WidgetActionsApi {
entityId?: EntityId, entityName?: string, additionalParams?: any, entityLabel?: string) => void;
elementClick: ($event: Event) => void;
getActiveEntityInfo: () => SubscriptionEntityInfo;
+ openDashboardStateInSeparateDialog: (targetDashboardStateId: string, params?: StateParams, dialogTitle?: string,
+ hideDashboardToolbar?: boolean, dialogWidth?: number, dialogHeight?: number) => void;
}
export interface AliasInfo {
diff --git a/ui-ngx/src/app/core/http/attribute.service.ts b/ui-ngx/src/app/core/http/attribute.service.ts
index 1bed864ab1..df7941cef3 100644
--- a/ui-ngx/src/app/core/http/attribute.service.ts
+++ b/ui-ngx/src/app/core/http/attribute.service.ts
@@ -50,11 +50,17 @@ export class AttributeService {
}
public deleteEntityTimeseries(entityId: EntityId, timeseries: Array, deleteAllDataForKeys = false,
- config?: RequestConfig): Observable {
+ startTs?: number, endTs?: number, config?: RequestConfig): Observable {
const keys = timeseries.map(attribute => encodeURI(attribute.key)).join(',');
- return this.http.delete(`/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/timeseries/delete` +
- `?keys=${keys}&deleteAllDataForKeys=${deleteAllDataForKeys}`,
- defaultHttpOptionsFromConfig(config));
+ let url = `/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/timeseries/delete` +
+ `?keys=${keys}&deleteAllDataForKeys=${deleteAllDataForKeys}`;
+ if (isDefinedAndNotNull(startTs)) {
+ url += `&startTs=${startTs}`;
+ }
+ if (isDefinedAndNotNull(endTs)) {
+ url += `&endTs=${endTs}`;
+ }
+ return this.http.delete(url, defaultHttpOptionsFromConfig(config));
}
public saveEntityAttributes(entityId: EntityId, attributeScope: AttributeScope, attributes: Array,
@@ -97,7 +103,7 @@ export class AttributeService {
});
let deleteEntityTimeseriesObservable: Observable;
if (deleteTimeseries.length) {
- deleteEntityTimeseriesObservable = this.deleteEntityTimeseries(entityId, deleteTimeseries, true, config);
+ deleteEntityTimeseriesObservable = this.deleteEntityTimeseries(entityId, deleteTimeseries, true, null, null, config);
} else {
deleteEntityTimeseriesObservable = of(null);
}
diff --git a/ui-ngx/src/app/modules/home/components/entity/contact-based.component.ts b/ui-ngx/src/app/modules/home/components/entity/contact-based.component.ts
index 00e94b9889..98c856eb9d 100644
--- a/ui-ngx/src/app/modules/home/components/entity/contact-based.component.ts
+++ b/ui-ngx/src/app/modules/home/components/entity/contact-based.component.ts
@@ -18,7 +18,7 @@ import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { FormBuilder, FormGroup, ValidatorFn, Validators } from '@angular/forms';
import { ContactBased } from '@shared/models/contact-based.model';
-import { AfterViewInit, Directive } from '@angular/core';
+import { AfterViewInit, ChangeDetectorRef, Directive } from '@angular/core';
import { POSTAL_CODE_PATTERNS } from '@home/models/contact.models';
import { HasId } from '@shared/models/base-data';
import { EntityComponent } from './entity.component';
@@ -30,8 +30,9 @@ export abstract class ContactBasedComponent> exten
protected constructor(protected store: Store,
protected fb: FormBuilder,
protected entityValue: T,
- protected entitiesTableConfigValue: EntityTableConfig) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ protected entitiesTableConfigValue: EntityTableConfig,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
buildForm(entity: T): FormGroup {
diff --git a/ui-ngx/src/app/modules/home/components/entity/entity.component.ts b/ui-ngx/src/app/modules/home/components/entity/entity.component.ts
index f5d53056a6..8d868eff98 100644
--- a/ui-ngx/src/app/modules/home/components/entity/entity.component.ts
+++ b/ui-ngx/src/app/modules/home/components/entity/entity.component.ts
@@ -17,7 +17,7 @@
import { BaseData, HasId } from '@shared/models/base-data';
import { FormBuilder, FormGroup } from '@angular/forms';
import { PageComponent } from '@shared/components/page.component';
-import { Directive, EventEmitter, Input, OnInit, Output } from '@angular/core';
+import { ChangeDetectorRef, Directive, EventEmitter, Input, OnInit, Output } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityAction } from '@home/models/entity/entity-component.models';
@@ -50,6 +50,7 @@ export abstract class EntityComponent,
@Input()
set isEdit(isEdit: boolean) {
this.isEditValue = isEdit;
+ this.cd.markForCheck();
this.updateFormState();
}
@@ -80,7 +81,8 @@ export abstract class EntityComponent,
protected constructor(protected store: Store,
protected fb: FormBuilder,
protected entityValue: T,
- protected entitiesTableConfigValue: C) {
+ protected entitiesTableConfigValue: C,
+ protected cd: ChangeDetectorRef) {
super(store);
this.entityForm = this.buildForm(this.entityValue);
}
diff --git a/ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts b/ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts
index efaa234adb..498b445440 100644
--- a/ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts
+++ b/ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject, Input, Optional } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject, Input, Optional } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { FormBuilder, FormGroup, Validators } from '@angular/forms';
@@ -77,8 +77,9 @@ export class DeviceProfileComponent extends EntityComponent {
protected translate: TranslateService,
@Optional() @Inject('entity') protected entityValue: DeviceProfile,
@Optional() @Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- protected fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ protected fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts
index afdb3eafc1..f8685cc193 100644
--- a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts
+++ b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject, Input, Optional } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject, Input, Optional } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { FormBuilder, FormGroup, Validators } from '@angular/forms';
@@ -43,8 +43,9 @@ export class TenantProfileComponent extends EntityComponent {
protected translate: TranslateService,
@Optional() @Inject('entity') protected entityValue: TenantProfile,
@Optional() @Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- protected fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ protected fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts
index dee17a000c..fc7637baef 100644
--- a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts
+++ b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts
@@ -284,7 +284,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
getActionDescriptors: this.getActionDescriptors.bind(this),
handleWidgetAction: this.handleWidgetAction.bind(this),
elementClick: this.elementClick.bind(this),
- getActiveEntityInfo: this.getActiveEntityInfo.bind(this)
+ getActiveEntityInfo: this.getActiveEntityInfo.bind(this),
+ openDashboardStateInSeparateDialog: this.openDashboardStateInSeparateDialog.bind(this)
};
this.widgetContext.customHeaderActions = [];
@@ -1025,7 +1026,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
this.updateEntityParams(params, targetEntityParamName, targetEntityId, entityName, entityLabel);
if (type === WidgetActionType.openDashboardState) {
if (descriptor.openInSeparateDialog) {
- this.openDashboardStateInDialog(descriptor, entityId, entityName, additionalParams, entityLabel);
+ this.openDashboardStateInSeparateDialog(descriptor.targetDashboardStateId, params, descriptor.dialogTitle,
+ descriptor.dialogHideDashboardToolbar, descriptor.dialogWidth, descriptor.dialogHeight);
} else {
this.widgetContext.stateController.openState(targetDashboardStateId, params, descriptor.openRightLayout);
}
@@ -1276,22 +1278,15 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
}
}
- private openDashboardStateInDialog(descriptor: WidgetActionDescriptor,
- entityId?: EntityId, entityName?: string, additionalParams?: any, entityLabel?: string) {
+ private openDashboardStateInSeparateDialog(targetDashboardStateId: string, params?: StateParams, dialogTitle?: string,
+ hideDashboardToolbar = true, dialogWidth?: number, dialogHeight?: number) {
const dashboard = deepClone(this.widgetContext.stateController.dashboardCtrl.dashboardCtx.getDashboard());
const stateObject: StateObject = {};
- stateObject.params = {};
- const targetEntityParamName = descriptor.stateEntityParamName;
- const targetDashboardStateId = descriptor.targetDashboardStateId;
- let targetEntityId: EntityId;
- if (descriptor.setEntityId) {
- targetEntityId = entityId;
- }
- this.updateEntityParams(stateObject.params, targetEntityParamName, targetEntityId, entityName, entityLabel);
+ stateObject.params = params;
if (targetDashboardStateId) {
stateObject.id = targetDashboardStateId;
}
- let title = descriptor.dialogTitle;
+ let title = dialogTitle;
if (!title) {
if (targetDashboardStateId && dashboard.configuration.states) {
const dashboardState = dashboard.configuration.states[targetDashboardStateId];
@@ -1304,7 +1299,6 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
title = dashboard.title;
}
title = this.utils.customTranslation(title, title);
- const params = stateObject.params;
const paramsEntityName = params && params.entityName ? params.entityName : '';
const paramsEntityLabel = params && params.entityLabel ? params.entityLabel : '';
title = insertVariable(title, 'entityName', paramsEntityName);
@@ -1324,28 +1318,27 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI
dashboard,
state: objToBase64([ stateObject ]),
title,
- hideToolbar: descriptor.dialogHideDashboardToolbar,
- width: descriptor.dialogWidth,
- height: descriptor.dialogHeight
+ hideToolbar: hideDashboardToolbar,
+ width: dialogWidth,
+ height: dialogHeight
}
});
}
private elementClick($event: Event) {
- const e = ($event.target || $event.srcElement) as Element;
- if (e.id) {
- const descriptors = this.getActionDescriptors('elementClick');
- if (descriptors.length) {
- descriptors.forEach((descriptor) => {
- if (descriptor.name === e.id) {
- $event.stopPropagation();
- const entityInfo = this.getActiveEntityInfo();
- const entityId = entityInfo ? entityInfo.entityId : null;
- const entityName = entityInfo ? entityInfo.entityName : null;
- const entityLabel = entityInfo && entityInfo.entityLabel ? entityInfo.entityLabel : null;
- this.handleWidgetAction($event, descriptor, entityId, entityName, null, entityLabel);
- }
- });
+ const elementClicked = ($event.target || $event.srcElement) as Element;
+ const descriptors = this.getActionDescriptors('elementClick');
+ if (descriptors.length) {
+ const idsList = descriptors.map(descriptor => `#${descriptor.name}`).join(',');
+ const targetElement = $(elementClicked).closest(idsList, this.widgetContext.$container[0]);
+ if (targetElement.length && targetElement[0].id) {
+ $event.stopPropagation();
+ const descriptor = descriptors.find(descriptorInfo => descriptorInfo.name === targetElement[0].id);
+ const entityInfo = this.getActiveEntityInfo();
+ const entityId = entityInfo ? entityInfo.entityId : null;
+ const entityName = entityInfo ? entityInfo.entityName : null;
+ const entityLabel = entityInfo && entityInfo.entityLabel ? entityInfo.entityLabel : null;
+ this.handleWidgetAction($event, descriptor, entityId, entityName, null, entityLabel);
}
}
}
diff --git a/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html b/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html
index 1c66d02275..3a40849616 100644
--- a/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html
+++ b/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.html
@@ -39,7 +39,7 @@
admin.smtp-protocol
-
+
{{protocol.toUpperCase()}}
@@ -127,7 +127,10 @@
-
+
+ {{ 'admin.change-password' | translate }}
+
+
common.password
diff --git a/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.ts b/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.ts
index 1d86318e8d..9750319565 100644
--- a/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/admin/mail-server.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, OnInit } from '@angular/core';
+import { Component, OnDestroy, OnInit } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { PageComponent } from '@shared/components/page.component';
@@ -25,21 +25,26 @@ import { AdminService } from '@core/http/admin.service';
import { ActionNotificationShow } from '@core/notification/notification.actions';
import { TranslateService } from '@ngx-translate/core';
import { HasConfirmForm } from '@core/guards/confirm-on-exit.guard';
-import { isString } from '@core/utils';
+import { isDefinedAndNotNull, isString } from '@core/utils';
+import { Subject } from 'rxjs';
+import { takeUntil } from 'rxjs/operators';
@Component({
selector: 'tb-mail-server',
templateUrl: './mail-server.component.html',
styleUrls: ['./mail-server.component.scss', './settings-card.scss']
})
-export class MailServerComponent extends PageComponent implements OnInit, HasConfirmForm {
+export class MailServerComponent extends PageComponent implements OnInit, OnDestroy, HasConfirmForm {
mailSettings: FormGroup;
adminSettings: AdminSettings;
smtpProtocols = ['smtp', 'smtps'];
+ showChangePassword = false;
tlsVersions = ['TLSv1', 'TLSv1.1', 'TLSv1.2', 'TLSv1.3'];
+ private destroy$ = new Subject();
+
constructor(protected store: Store,
private router: Router,
private adminService: AdminService,
@@ -56,12 +61,22 @@ export class MailServerComponent extends PageComponent implements OnInit, HasCon
if (this.adminSettings.jsonValue && isString(this.adminSettings.jsonValue.enableTls)) {
this.adminSettings.jsonValue.enableTls = (this.adminSettings.jsonValue.enableTls as any) === 'true';
}
+ this.showChangePassword =
+ isDefinedAndNotNull(this.adminSettings.jsonValue.showChangePassword) ? this.adminSettings.jsonValue.showChangePassword : true ;
+ delete this.adminSettings.jsonValue.showChangePassword;
this.mailSettings.reset(this.adminSettings.jsonValue);
+ this.enableMailPassword(!this.showChangePassword);
this.enableProxyChanged();
}
);
}
+ ngOnDestroy() {
+ this.destroy$.next();
+ this.destroy$.complete();
+ super.ngOnDestroy();
+ }
+
buildMailServerSettingsForm() {
this.mailSettings = this.fb.group({
mailFrom: ['', [Validators.required]],
@@ -81,14 +96,23 @@ export class MailServerComponent extends PageComponent implements OnInit, HasCon
proxyUser: [''],
proxyPassword: [''],
username: [''],
+ changePassword: [false],
password: ['']
});
this.registerDisableOnLoadFormControl(this.mailSettings.get('smtpProtocol'));
this.registerDisableOnLoadFormControl(this.mailSettings.get('enableTls'));
this.registerDisableOnLoadFormControl(this.mailSettings.get('enableProxy'));
- this.mailSettings.get('enableProxy').valueChanges.subscribe(() => {
+ this.registerDisableOnLoadFormControl(this.mailSettings.get('changePassword'));
+ this.mailSettings.get('enableProxy').valueChanges.pipe(
+ takeUntil(this.destroy$)
+ ).subscribe(() => {
this.enableProxyChanged();
});
+ this.mailSettings.get('changePassword').valueChanges.pipe(
+ takeUntil(this.destroy$)
+ ).subscribe((value) => {
+ this.enableMailPassword(value);
+ });
}
enableProxyChanged(): void {
@@ -102,8 +126,16 @@ export class MailServerComponent extends PageComponent implements OnInit, HasCon
}
}
+ enableMailPassword(enable: boolean) {
+ if (enable) {
+ this.mailSettings.get('password').enable({emitEvent: false});
+ } else {
+ this.mailSettings.get('password').disable({emitEvent: false});
+ }
+ }
+
sendTestMail(): void {
- this.adminSettings.jsonValue = {...this.adminSettings.jsonValue, ...this.mailSettings.value};
+ this.adminSettings.jsonValue = {...this.adminSettings.jsonValue, ...this.mailSettingsFormValue};
this.adminService.sendTestMail(this.adminSettings).subscribe(
() => {
this.store.dispatch(new ActionNotificationShow({ message: this.translate.instant('admin.test-mail-sent'),
@@ -113,13 +145,11 @@ export class MailServerComponent extends PageComponent implements OnInit, HasCon
}
save(): void {
- this.adminSettings.jsonValue = {...this.adminSettings.jsonValue, ...this.mailSettings.value};
+ this.adminSettings.jsonValue = {...this.adminSettings.jsonValue, ...this.mailSettingsFormValue};
this.adminService.saveAdminSettings(this.adminSettings).subscribe(
(adminSettings) => {
- if (!adminSettings.jsonValue.password) {
- adminSettings.jsonValue.password = this.mailSettings.value.password;
- }
this.adminSettings = adminSettings;
+ this.showChangePassword = true;
this.mailSettings.reset(this.adminSettings.jsonValue);
}
);
@@ -129,4 +159,9 @@ export class MailServerComponent extends PageComponent implements OnInit, HasCon
return this.mailSettings;
}
+ private get mailSettingsFormValue(): MailServerSettings {
+ const formValue = this.mailSettings.value;
+ delete formValue.changePassword;
+ return formValue;
+ }
}
diff --git a/ui-ngx/src/app/modules/home/pages/admin/resource/resources-library.component.ts b/ui-ngx/src/app/modules/home/pages/admin/resource/resources-library.component.ts
index 50d565ccc0..a6a14106d5 100644
--- a/ui-ngx/src/app/modules/home/pages/admin/resource/resources-library.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/admin/resource/resources-library.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject, OnDestroy, OnInit } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject, OnDestroy, OnInit } from '@angular/core';
import { Subject } from 'rxjs';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
@@ -30,7 +30,7 @@ import {
ResourceTypeTranslationMap
} from '@shared/models/resource.models';
import { pairwise, startWith, takeUntil } from 'rxjs/operators';
-import { ActionNotificationShow } from "@core/notification/notification.actions";
+import { ActionNotificationShow } from '@core/notification/notification.actions';
@Component({
selector: 'tb-resources-library',
@@ -48,8 +48,9 @@ export class ResourcesLibraryComponent extends EntityComponent impleme
protected translate: TranslateService,
@Inject('entity') protected entityValue: Resource,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
@@ -102,7 +103,7 @@ export class ResourcesLibraryComponent extends EntityComponent impleme
if (this.isAdd) {
form.addControl('data', this.fb.control(null, Validators.required));
}
- return form
+ return form;
}
updateForm(entity: Resource) {
diff --git a/ui-ngx/src/app/modules/home/pages/asset/asset.component.ts b/ui-ngx/src/app/modules/home/pages/asset/asset.component.ts
index 2ec441c2d3..0ce8f8e6e7 100644
--- a/ui-ngx/src/app/modules/home/pages/asset/asset.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/asset/asset.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -41,8 +41,9 @@ export class AssetComponent extends EntityComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: AssetInfo,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/customer/customer.component.ts b/ui-ngx/src/app/modules/home/pages/customer/customer.component.ts
index 20e49d118b..2c6e403767 100644
--- a/ui-ngx/src/app/modules/home/pages/customer/customer.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/customer/customer.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { FormBuilder, FormGroup, Validators } from '@angular/forms';
@@ -42,8 +42,9 @@ export class CustomerComponent extends ContactBasedComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: Customer,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- protected fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ protected fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/modules/home/pages/dashboard/dashboard-form.component.ts b/ui-ngx/src/app/modules/home/pages/dashboard/dashboard-form.component.ts
index 13a0340780..35e6daec22 100644
--- a/ui-ngx/src/app/modules/home/pages/dashboard/dashboard-form.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/dashboard/dashboard-form.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -49,8 +49,9 @@ export class DashboardFormComponent extends EntityComponent {
private dashboardService: DashboardService,
@Inject('entity') protected entityValue: Dashboard,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/device/device.component.ts b/ui-ngx/src/app/modules/home/pages/device/device.component.ts
index c82ad1c26a..b915d791ab 100644
--- a/ui-ngx/src/app/modules/home/pages/device/device.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/device/device.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -56,8 +56,9 @@ export class DeviceComponent extends EntityComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: DeviceInfo,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/edge/edge.component.ts b/ui-ngx/src/app/modules/home/pages/edge/edge.component.ts
index 007176ce24..b8539df134 100644
--- a/ui-ngx/src/app/modules/home/pages/edge/edge.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/edge/edge.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '@home/components/entity/entity.component';
@@ -42,8 +42,9 @@ export class EdgeComponent extends EntityComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: EdgeInfo,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/entity-view/entity-view.component.ts b/ui-ngx/src/app/modules/home/pages/entity-view/entity-view.component.ts
index 827e83a7be..38b547a012 100644
--- a/ui-ngx/src/app/modules/home/pages/entity-view/entity-view.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/entity-view/entity-view.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -53,8 +53,9 @@ export class EntityViewComponent extends EntityComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: EntityViewInfo,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/ota-update/ota-update.component.ts b/ui-ngx/src/app/modules/home/pages/ota-update/ota-update.component.ts
index 7ea9ef2e3b..e88f2b10a7 100644
--- a/ui-ngx/src/app/modules/home/pages/ota-update/ota-update.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/ota-update/ota-update.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject, OnDestroy, OnInit } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject, OnDestroy, OnInit } from '@angular/core';
import { combineLatest, Subject } from 'rxjs';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
@@ -50,8 +50,9 @@ export class OtaUpdateComponent extends EntityComponent implements O
protected translate: TranslateService,
@Inject('entity') protected entityValue: OtaPackage,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/rulechain/rulechain.component.ts b/ui-ngx/src/app/modules/home/pages/rulechain/rulechain.component.ts
index 3d8eae1e7a..9aa518e4c0 100644
--- a/ui-ngx/src/app/modules/home/pages/rulechain/rulechain.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/rulechain/rulechain.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -37,8 +37,9 @@ export class RuleChainComponent extends EntityComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: RuleChain,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
ngOnInit() {
diff --git a/ui-ngx/src/app/modules/home/pages/tenant/tenant.component.ts b/ui-ngx/src/app/modules/home/pages/tenant/tenant.component.ts
index 30c0dc36aa..c59a9d557b 100644
--- a/ui-ngx/src/app/modules/home/pages/tenant/tenant.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/tenant/tenant.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { FormBuilder, FormGroup, Validators } from '@angular/forms';
@@ -36,8 +36,9 @@ export class TenantComponent extends ContactBasedComponent {
protected translate: TranslateService,
@Inject('entity') protected entityValue: TenantInfo,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- protected fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ protected fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/modules/home/pages/user/user.component.ts b/ui-ngx/src/app/modules/home/pages/user/user.component.ts
index f6fc69b7d6..7187731b5c 100644
--- a/ui-ngx/src/app/modules/home/pages/user/user.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/user/user.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject, Optional } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject, Optional } from '@angular/core';
import { select, Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -43,8 +43,9 @@ export class UserComponent extends EntityComponent {
constructor(protected store: Store,
@Optional() @Inject('entity') protected entityValue: User,
@Optional() @Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/modules/home/pages/widget/widgets-bundle.component.ts b/ui-ngx/src/app/modules/home/pages/widget/widgets-bundle.component.ts
index 4947a90e1f..dd3b827c1c 100644
--- a/ui-ngx/src/app/modules/home/pages/widget/widgets-bundle.component.ts
+++ b/ui-ngx/src/app/modules/home/pages/widget/widgets-bundle.component.ts
@@ -14,7 +14,7 @@
/// limitations under the License.
///
-import { Component, Inject } from '@angular/core';
+import { ChangeDetectorRef, Component, Inject } from '@angular/core';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { EntityComponent } from '../../components/entity/entity.component';
@@ -32,8 +32,9 @@ export class WidgetsBundleComponent extends EntityComponent {
constructor(protected store: Store,
@Inject('entity') protected entityValue: WidgetsBundle,
@Inject('entitiesTableConfig') protected entitiesTableConfigValue: EntityTableConfig,
- public fb: FormBuilder) {
- super(store, fb, entityValue, entitiesTableConfigValue);
+ public fb: FormBuilder,
+ protected cd: ChangeDetectorRef) {
+ super(store, fb, entityValue, entitiesTableConfigValue, cd);
}
hideDelete() {
diff --git a/ui-ngx/src/app/shared/models/settings.models.ts b/ui-ngx/src/app/shared/models/settings.models.ts
index 3877488671..8b2d4aac78 100644
--- a/ui-ngx/src/app/shared/models/settings.models.ts
+++ b/ui-ngx/src/app/shared/models/settings.models.ts
@@ -27,6 +27,7 @@ export interface AdminSettings {
export declare type SmtpProtocol = 'smtp' | 'smtps';
export interface MailServerSettings {
+ showChangePassword: boolean;
mailFrom: string;
smtpProtocol: SmtpProtocol;
smtpHost: string;
@@ -34,7 +35,8 @@ export interface MailServerSettings {
timeout: number;
enableTls: boolean;
username: string;
- password: string;
+ changePassword?: boolean;
+ password?: string;
enableProxy: boolean;
proxyHost: string;
proxyPort: number;
diff --git a/ui-ngx/src/assets/locale/locale.constant-de_DE.json b/ui-ngx/src/assets/locale/locale.constant-de_DE.json
index 32aa7272a1..eee8d89fd7 100644
--- a/ui-ngx/src/assets/locale/locale.constant-de_DE.json
+++ b/ui-ngx/src/assets/locale/locale.constant-de_DE.json
@@ -803,7 +803,7 @@
"rulechain-templates": "Regelkettenvorlagen",
"rulechains": "Rand Regelketten",
"search": "Kanten durchsuchen",
- "selected-edges": "{Anzahl, Plural, 1 {1 Kante} andere {# Kanten}} ausgewählt",
+ "selected-edges": "{count, plural, 1 {1 Rand} other {# Rand} } ausgewählt",
"any-edge": "Beliebige Kante",
"no-edge-types-matching": "Es wurden keine Kantentypen gefunden, die mit '{{entitySubtype}}' übereinstimmen.",
"edge-type-list-empty": "Keine Kantentypen ausgewählt.",
@@ -1452,7 +1452,7 @@
"unset-auto-assign-to-edge-text": "Nach der Bestätigung wird die Kantenregelkette bei der Erstellung nicht mehr automatisch den Kanten zugewiesen.",
"edge-template-root": "Vorlagenstamm",
"search": "Suchen Sie nach Regelketten",
- "selected-rulechains": "{count, plural, 1 {1 Regelkette} andere {# Regelketten}} ausgewählt",
+ "selected-rulechains": "{count, plural, 1 {1 Regelkette} other {# Regelketten} } ausgewählt",
"open-rulechain": "Regelkette öffnen",
"assign-to-edge": "Rand zuweisen",
"edge-rulechain": "Kantenregelkette"
diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json
index eceef34751..8d1632a241 100644
--- a/ui-ngx/src/assets/locale/locale.constant-en_US.json
+++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json
@@ -104,6 +104,7 @@
"proxy-port-range": "Proxy port should be in a range from 1 to 65535.",
"proxy-user": "Proxy user",
"proxy-password": "Proxy password",
+ "change-password": "Change password",
"send-test-mail": "Send test mail",
"sms-provider": "SMS provider",
"sms-provider-settings": "SMS provider settings",
diff --git a/ui-ngx/src/assets/locale/locale.constant-es_ES.json b/ui-ngx/src/assets/locale/locale.constant-es_ES.json
index 9971ab479a..a782298bdc 100644
--- a/ui-ngx/src/assets/locale/locale.constant-es_ES.json
+++ b/ui-ngx/src/assets/locale/locale.constant-es_ES.json
@@ -413,8 +413,8 @@
"unassign-asset-from-edge": "Anular activo de bodre",
"unassign-asset-from-edge-title": "¿Está seguro de que desea desasignar el activo '{{assetName}}'?",
"unassign-asset-from-edge-text": "Después de la confirmación, el activo no será asignado y el borde no podrá acceder a él",
- "unassign-assets-from-edge-action-title": "Anular asignación {count, plural, 1 {1 activo} other {# activos}} desde el borde",
- "unassign-assets-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 activo} other {# activos}}?",
+ "unassign-assets-from-edge-action-title": "Anular asignación {count, plural, 1 {1 activo} other {# activos} } desde el borde",
+ "unassign-assets-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 activo} other {# activos} }?",
"unassign-assets-from-edge-text": "Después de la confirmación, todos los activos seleccionados quedarán sin asignar y el borde no podrá acceder a ellos."
},
"attribute": {
@@ -950,7 +950,7 @@
"assign-device-to-edge-text": "Seleccione los dispositivos para asignar al borde",
"unassign-device-from-edge-title": "¿Está seguro de que desea desasignar el dispositivo '{{deviceName}}'?",
"unassign-device-from-edge-text": "Después de la confirmación, el dispositivo no será asignado y el borde no podrá acceder a él",
- "unassign-devices-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 dispositivo} other {# dispositivos}}?",
+ "unassign-devices-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 dispositivo} other {# dispositivos} }?",
"unassign-devices-from-edge-text": "Después de la confirmación, todos los dispositivos seleccionados quedarán sin asignar y el borde no podrá acceder a ellos."
},
"device-profile": {
@@ -1123,7 +1123,7 @@
"delete": "Eliminar borde",
"delete-edge-title": "¿Está seguro de que desea eliminar el borde '{{edgeName}}'?",
"delete-edge-text": "Tenga cuidado, después de la confirmación, el borde y todos los datos relacionados serán irrecuperables",
- "delete-edges-title": "¿Está seguro de que desea edge {count, plural, 1 {1 borde} other {# bordes}}?",
+ "delete-edges-title": "¿Está seguro de que desea edge {count, plural, 1 {1 borde} other {# bordes} }?",
"delete-edges-text": "Tenga cuidado, después de la confirmación se eliminarán todos los bordes seleccionados y todos los datos relacionados se volverán irrecuperables",
"name": "Nombre",
"name-starts-with": "Edge name starts with",
@@ -1156,7 +1156,7 @@
"unassign-from-customer": "Anular asignación del cliente",
"unassign-edge-title": "¿Está seguro de que desea desasignar el borde '{{edgeName}}'?",
"unassign-edge-text": "Después de la confirmación, el borde quedará sin asignar y el cliente no podrá acceder a él",
- "unassign-edges-title": "¿Está seguro de que desea anular la asignación de {count, plural, 1 {1 borde} other {# bordes}}?",
+ "unassign-edges-title": "¿Está seguro de que desea anular la asignación de {count, plural, 1 {1 borde} other {# bordes} }?",
"unassign-edges-text": "Después de la confirmación de todos los bordes seleccionados, se anulará la asignación y el cliente no podrá acceder a ellos.",
"make-public": "Hacer público el borde",
"make-public-edge-title": "¿Estás seguro de que quieres hacer público el edge '{{edgeName}}'?",
@@ -1189,14 +1189,14 @@
"rulechain-templates": "Plantillas, de cadena de reglas",
"rulechains": "Cadenas de regla de borde",
"search": "Bordes de búsqueda",
- "selected-edges": "{count, plural, 1 {1 borde} other {# bordes}} seleccionados",
+ "selected-edges": "{count, plural, 1 {1 borde} other {# bordes} } seleccionadas",
"any-edge": "Cualquier bordee",
"no-edge-types-matching": "No se encontraron tipos de aristas que coincidan con '{{entitySubtype}}'.",
"edge-type-list-empty": "No se seleccionó ningún tipo de borde.",
"edge-types": "Tipos de bordes",
"enter-edge-type": "Ingrese el tipo de borde",
"deployed": "Desplegada",
- "pending": "Pending",
+ "pending": "Pendiente",
"downlinks": "Enlaces descendentes",
"no-downlinks-prompt": "No se encontraron enlaces descendentes",
"sync-process-started-successfully": "¡El proceso de sincronización se inició correctamente!",
@@ -1356,7 +1356,7 @@
"type-api-usage-state": "Estado de uso de la API",
"type-edge": "Borde",
"type-edges": "Bordes",
- "list-of-edges": "{cuenta, plural, 1 {Un borde} other {Lista de # bordes}}",
+ "list-of-edges": "{count, plural, 1 {Un borde} other {Lista de # bordes} }",
"edge-name-starts-with": "Bordes cuyos nombres comienzan con '{{prefijo}}'"
},
"entity-field": {
@@ -1481,9 +1481,9 @@
"assign-entity-view-to-edge-text": "Seleccione las vistas de entidad para asignar al borde",
"unassign-entity-view-from-edge-title": "¿Está seguro de que desea anular la asignación de la vista de entidad '{{entityViewName}}'?",
"unassign-entity-view-from-edge-text": "Después de la confirmación, la vista de entidad quedará sin asignar y el borde no podrá acceder a ella",
- "unassign-entity-views-from-edge-action-title": "Anular asignación {recuento, plural, 1 {1 vista de entidad} otras {# vistas de entidad}} del borde",
+ "unassign-entity-views-from-edge-action-title": "Anular asignación {count, plural, 1 {1 vista de entidad} other {# vistas de entidad} } del borde",
"unassign-entity-view-from-edge": "Anular asignación de vista de entidad",
- "unassign-entity-views-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 vista de entidad} other {# vistas de entidad}}?",
+ "unassign-entity-views-from-edge-title": "¿Está seguro de que desea desasignar {count, plural, 1 {1 vista de entidad} other {# vistas de entidad} }?",
"unassign-entity-views-from-edge-text": "Después de la confirmación, todas las vistas de entidad seleccionadas no serán asignadas y el borde no podrá acceder a ellas"
},
"event": {
@@ -2074,9 +2074,9 @@
"delete-rulechains": "Eliminar cadenas de reglas",
"unassign-rulechain": "Anular asignación de cadena de reglas",
"unassign-rulechains": "Anular asignación de cadenas de reglas",
- "unassign-rulechain-title": "¿Está seguro de que desea desasignar la cadena de reglas '{{ruleChainTitle}}'?",
+ "unassign-rulechain-title": "¿Está seguro de que desea desasignar la cadena de reglas '{{ruleChainName}}'?",
"unassign-rulechain-from-edge-text": "Después de la confirmación, la cadena de reglas quedará sin asignar y el borde no podrá acceder a ella",
- "unassign-rulechains-from-edge-action-title": "Anular asignación {count, plural, 1 {1 cadena de reglas} other {# cadenas de reglas}} des bordes",
+ "unassign-rulechains-from-edge-action-title": "Anular asignación {count, plural, 1 {1 cadena de reglas} other {# cadenas de reglas} } des bordes",
"unassign-rulechains-from-edge-text": "Después de la confirmación, todas las cadenas de reglas seleccionadas quedarán sin asignar y el borde no podrá acceder a ellas",
"assign-rulechain-to-edge-title": "Asignar cadena (s) de reglas a borde",
"assign-rulechain-to-edge-text": "Seleccione las cadenas de reglas para asignar al borde",
diff --git a/ui-ngx/src/assets/locale/locale.constant-fr_FR.json b/ui-ngx/src/assets/locale/locale.constant-fr_FR.json
index 45cc24c5cd..facf0f8534 100644
--- a/ui-ngx/src/assets/locale/locale.constant-fr_FR.json
+++ b/ui-ngx/src/assets/locale/locale.constant-fr_FR.json
@@ -736,7 +736,7 @@
"assign-device-to-edge-text":"Veuillez sélectionner la bordure pour attribuer le ou les dispositifs",
"unassign-device-from-edge-title": "Êtes-vous sûr de vouloir annuler l'affection du dispositif {{deviceName}} '?",
"unassign-device-from-edge-text": "Après la confirmation, dispositif sera non attribué et ne sera pas accessible a la bordure.",
- "unassign-devices-from-edge-title": "Voulez-vous vraiment annuler l'affectation de {count, plural, 1 {1 device} other {# devices}}?",
+ "unassign-devices-from-edge-title": "Voulez-vous vraiment annuler l'affectation de {count, plural, 1 {1 device} other {# devices} }?",
"unassign-devices-from-edge-text": "Après la confirmation, tous les dispositifs sélectionnés ne seront pas attribues et ne seront pas accessibles par la bordure."
},
"dialog": {
@@ -755,7 +755,7 @@
"delete": "Supprimer la bordure",
"delete-edge-title": "Êtes-vous sûr de vouloir supprimer la bordure '{{edgeName}}'?",
"delete-edge-text": "Faites attention, après la confirmation, la bordure et toutes les données associées deviendront irrécupérables",
- "delete-edges-title": "Êtes-vous sûr de vouloir supprimer {count, plural, 1 {1 bordure} other {# bordure}}?",
+ "delete-edges-title": "Êtes-vous sûr de vouloir supprimer {count, plural, 1 {1 bordure} other {# bordure} }?",
"delete-edges-text": "Faites attention, après la confirmation, tous les bordures sélectionnés seront supprimés et toutes les données associées deviendront irrécupérables.",
"name": "Nom",
"name-starts-with": "Le nom du bord commence par",
@@ -788,7 +788,7 @@
"unassign-from-customer": "Retirer du client",
"unassign-edge-title": "Êtes-vous sûr de vouloir annuler l'affection du dispositif {{edgeName}}",
"unassign-edge-text": "Après la confirmation, le dispositif ne sera pas attribué et ne sera pas accessible au client",
- "unassign-edges-title": "Voulez-vous vraiment annuler l'attribution de {count, plural, 1 {1 bordure} other {# bordures}}?",
+ "unassign-edges-title": "Voulez-vous vraiment annuler l'attribution de {count, plural, 1 {1 bordure} other {# bordures} }?",
"unassign-edges-text": "Après la confirmation, tous les bordures sélectionnés ne seront plus attribués et ne seront pas accessibles par le client.",
"make-public": "Make edge public",
"make-public-edge-title": "Are you sure you want to make the edge '{{edgeName}}' public?",
@@ -821,7 +821,7 @@
"rulechain-templates": "Modèles de chaîne de règles",
"rulechains": "Chaînes de règles de la bordure",
"search": "Rechercher les bords",
- "selected-edges": "{count, plural, 1 {1 edge} other {# bords}} sélectionné",
+ "selected-edges": "{count, plural, 1 {1 bordure} other {# bords} } sélectionné",
"any-edge": "Tout bord",
"no-edge-types-matching": "Aucun type d'arête correspondant à \"{{entitySubtype}}\" n'a été trouvé.",
"edge-type-list-empty": "Aucun type d'arête sélectionné.",
@@ -1479,9 +1479,9 @@
"delete-rulechains": "Supprimer une chaînes de règles",
"unassign-rulechain": "Retirer chaîne de règles",
"unassign-rulechains": "Retirer chaînes de règles",
- "unassign-rulechain-title": "AÊtes-vous sûr de vouloir retirer l'attribution de chaînes de règles '{{ruleChainTitle}}'?",
+ "unassign-rulechain-title": "AÊtes-vous sûr de vouloir retirer l'attribution de chaînes de règles '{{ruleChainName}}'?",
"unassign-rulechain-from-edge-text": "Après la confirmation, l'actif sera non attribué et ne sera pas accessible a la bordure.",
- "unassign-rulechains-from-edge-action-title": "Retirer {count, plural, 1 {1 chaîne de règles} other {# chaînes de règles}} de la bordure",
+ "unassign-rulechains-from-edge-action-title": "Retirer {count, plural, 1 {1 chaîne de règles} other {# chaînes de règles} } de la bordure",
"unassign-rulechains-from-edge-text": "Après la confirmation, tous les chaînes de règles sélectionnés ne seront pas attribués et ne seront pas accessibles a la bordure.",
"assign-rulechain-to-edge-title": "Attribuer les chaînes de règles a la bordure",
"assign-rulechain-to-edge-text": "Veuillez sélectionner la bordure pour attribuer le ou les chaînes de règles",
@@ -1497,11 +1497,11 @@
"unset-auto-assign-to-edge-text": "Après la confirmation, la chaîne de règles d'arêtes ne sera plus automatiquement affectée aux arêtes lors de la création.",
"edge-template-root": "Racine du modèle",
"search": "Rechercher des chaînes de règles",
- "selected-rulechains": "{count, plural, 1 {1 rule chain} other {# rule chains}} sélectionné",
+ "selected-rulechains": "{count, plural, 1 {1 rule chain} other {# rule chains} } sélectionné",
"open-rulechain": "Chaîne de règles ouverte",
- "assign-to-edge": "Attribuer à Edge",
- "edge-rulechain": "Chaîne de règles Edge",
- "unassign-rulechains-from-edge-title": "Voulez-vous vraiment annuler l'attribution de {count, plural, 1 {1 rulechain} other {# rulechains}}?"
+ "assign-to-edge": "Attribuer à Bordure",
+ "edge-rulechain": "Chaîne de règles Bordure",
+ "unassign-rulechains-from-edge-title": "Voulez-vous vraiment annuler l'attribution de {count, plural, 1 {1 rulechain} other {# rulechains} }?"
},
"rulenode": {
"add": "Ajouter un noeud de règle",