getLatestFromCacheOrFetchFromDb(TbContext ctx, TbMsg msg) {
+ EntityId originator = msg.getOriginator();
+ ValueWithTs valueWithTs = cache.get(msg.getOriginator());
+ return valueWithTs != null ? Futures.immediateFuture(valueWithTs) : fetchLatestValueAsync(ctx, originator);
+ }
- private ValueWithTs(long ts, double value) {
- this.ts = ts;
- this.value = value;
- }
+ private record ValueWithTs(long ts, double value) {
}
}
diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/SemaphoreWithTbMsgQueue.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/SemaphoreWithTbMsgQueue.java
new file mode 100644
index 0000000000..fa00856b4b
--- /dev/null
+++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/SemaphoreWithTbMsgQueue.java
@@ -0,0 +1,134 @@
+/**
+ * Copyright © 2016-2024 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.rule.engine.util;
+
+import com.google.common.util.concurrent.ListenableFuture;
+import lombok.Data;
+import lombok.extern.slf4j.Slf4j;
+import org.thingsboard.common.util.DonAsynchron;
+import org.thingsboard.rule.engine.api.TbContext;
+import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.msg.TbMsg;
+
+import java.util.Queue;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.Semaphore;
+import java.util.function.BiFunction;
+
+/**
+ * A utility class designed to manage a queue of messages for a specific entity, ensuring that
+ * message processing is synchronized on a per-entity basis. This is achieved through the use of a semaphore,
+ * allowing only one message at a time to be processed for each entity ID, thus preventing race conditions
+ * and ensuring thread-safe operations.
+ *
+ * This class is especially useful in scenarios where the order of message processing and
+ * resource access synchronization are crucial, such as updating caches or databases in a concurrent environment.
+ */
+@Data
+@Slf4j
+public class SemaphoreWithTbMsgQueue {
+
+ private final EntityId entityId;
+ private final Semaphore semaphore = new Semaphore(1);
+ private final Queue queue = new ConcurrentLinkedQueue<>();
+
+ /**
+ * Adds a message to the queue for asynchronous processing and attempts to process the queue if possible.
+ * This method is thread-safe and ensures that messages are processed in the order they were added,
+ * with each message for a specific entity being processed one at a time due to the semaphore control.
+ *
+ * @param msg The message to be processed.
+ * @param ctx The context in which the message should be processed.
+ * @param msgProcessingFunction The function that defines how the message will be processed.
+ */
+ public void addToQueueAndTryProcess(TbMsg msg, TbContext ctx, BiFunction> msgProcessingFunction) {
+ queue.add(new TbMsgTbContextBiFunction(msg, ctx, msgProcessingFunction));
+ tryProcessQueue();
+ }
+
+ /**
+ * Attempts to process the next message in the queue. If the semaphore is available (indicating
+ * that no other message for the same entity is currently being processed), this method will
+ * acquire the semaphore and start processing the message. If the semaphore is not available,
+ * this method will return immediately, ensuring that messages are processed sequentially
+ * for each entity.
+ *
+ * This method is automatically called after adding a message to the queue to ensure
+ * that the queue is processed promptly.
+ */
+ private void tryProcessQueue() {
+ while (!queue.isEmpty()) {
+ // The semaphore have to be acquired before EACH poll and released before NEXT poll.
+ // Otherwise, some message will remain unprocessed in queue
+ if (!semaphore.tryAcquire()) {
+ return;
+ }
+ TbMsgTbContextBiFunction tbMsgTbContext = null;
+ try {
+ tbMsgTbContext = queue.poll();
+ if (tbMsgTbContext == null) {
+ semaphore.release();
+ continue;
+ }
+ final TbMsg msg = tbMsgTbContext.msg();
+ if (!msg.getCallback().isMsgValid()) {
+ log.trace("[{}] Skipping non-valid message [{}]", entityId, msg);
+ semaphore.release();
+ continue;
+ }
+ //DO PROCESSING
+ final TbContext ctx = tbMsgTbContext.ctx();
+ final ListenableFuture resultMsgFuture = tbMsgTbContext.biFunction().apply(ctx, msg);
+ DonAsynchron.withCallback(resultMsgFuture, resultMsg -> {
+ try {
+ ctx.tellSuccess(resultMsg);
+ } finally {
+ semaphore.release();
+ tryProcessQueue();
+ }
+ }, t -> {
+ try {
+ ctx.tellFailure(msg, t);
+ } finally {
+ semaphore.release();
+ tryProcessQueue();
+ }
+ }, ctx.getDbCallbackExecutor());
+ } catch (Throwable t) {
+ semaphore.release();
+ if (tbMsgTbContext == null) { // if no message polled, the loop become infinite, will throw exception
+ log.error("[{}] Failed to process TbMsgTbContext queue", entityId, t);
+ throw t;
+ }
+ TbMsg msg = tbMsgTbContext.msg();
+ TbContext ctx = tbMsgTbContext.ctx();
+ log.warn("[{}] Failed to process message: {}", entityId, msg, t);
+ ctx.tellFailure(msg, t); // you are not allowed to throw here, because queue will remain unprocessed
+ continue; // We are probably the last who process the queue. We have to continue poll until get successful callback or queue is empty
+ }
+ break; //submitted async exact one task. next poll will try on callback
+ }
+ }
+
+ /**
+ * A utility record to hold the tuple of a {@link TbMsg}, {@link TbContext}, and the message processing function.
+ * This facilitates passing these three elements as a single object within the queue.
+ */
+ private record TbMsgTbContextBiFunction(TbMsg msg, TbContext ctx,
+ BiFunction> biFunction) {
+ }
+
+}
diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java
index e036e6333e..80e669847c 100644
--- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java
+++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java
@@ -20,7 +20,6 @@ import com.google.common.util.concurrent.Futures;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.tuple.Triple;
import org.assertj.core.api.SoftAssertions;
-import org.junit.Assert;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -69,7 +68,6 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
-import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyDouble;
@@ -78,9 +76,7 @@ import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willAnswer;
import static org.mockito.BDDMockito.willReturn;
-import static org.mockito.BDDMockito.willReturn;
import static org.mockito.BDDMockito.willThrow;
-import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
@@ -543,7 +539,7 @@ public class TbMathNodeTest {
ArgumentCaptor tCaptor = ArgumentCaptor.forClass(Throwable.class);
Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(eq(msg), tCaptor.capture());
- Assert.assertNotNull(tCaptor.getValue().getMessage());
+ assertNotNull(tCaptor.getValue().getMessage());
}
@Test
@@ -558,7 +554,7 @@ public class TbMathNodeTest {
ArgumentCaptor tCaptor = ArgumentCaptor.forClass(Throwable.class);
Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(eq(msg), tCaptor.capture());
- Assert.assertNotNull(tCaptor.getValue().getMessage());
+ assertNotNull(tCaptor.getValue().getMessage());
}
@Test
@@ -574,10 +570,10 @@ public class TbMathNodeTest {
List slowMsgList = IntStream.range(0, 5)
.mapToObj(x -> TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originatorSlow, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 2).put("b", 2).toString()))
- .collect(Collectors.toList());
+ .toList();
List fastMsgList = IntStream.range(0, 2)
.mapToObj(x -> TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originatorFast, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 2).put("b", 2).toString()))
- .collect(Collectors.toList());
+ .toList();
assertThat(slowMsgList.size()).as("slow msgs >= rule-dispatcher pool size").isGreaterThanOrEqualTo(RULE_DISPATCHER_POOL_SIZE);
@@ -714,7 +710,7 @@ public class TbMathNodeTest {
}).given(node).onMsg(any(), any());
return Triple.of(ctx, resultKey, node);
})
- .collect(Collectors.toList());
+ .toList();
ctxNodes.forEach(ctxNode -> ruleEngineDispatcherExecutor.executeAsync(() -> ctxNode.getRight()
.onMsg(ctxNode.getLeft(), TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, "{\"a\":2,\"b\":2}"))));
ctxNodes.forEach(ctxNode -> verify(ctxNode.getRight(), timeout(5000)).onMsg(eq(ctxNode.getLeft()), any()));
diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java
index 74269b9be6..c427dcbd3b 100644
--- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java
+++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java
@@ -16,20 +16,20 @@
package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures;
-import lombok.Data;
-import lombok.RequiredArgsConstructor;
-import org.assertj.core.api.Assertions;
+import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
+import org.junit.jupiter.params.provider.NullAndEmptySource;
+import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
-import org.mockito.ArgumentMatcher;
import org.mockito.Mock;
import org.mockito.Spy;
import org.mockito.junit.jupiter.MockitoExtension;
+import org.thingsboard.common.util.AbstractListeningExecutor;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest;
@@ -54,34 +54,48 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import java.util.List;
+import java.util.Optional;
import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
+import java.util.stream.IntStream;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.ArgumentMatchers.anySet;
import static org.mockito.ArgumentMatchers.anyString;
-import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.BDDMockito.willAnswer;
+import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.reset;
+import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;
+@Slf4j
@ExtendWith(MockitoExtension.class)
public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
- private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID());
- private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID());
- private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
+ private final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.fromString("2ba3ded4-882b-40cf-999a-89da9ccd58f9"));
+ private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("3842e740-0d89-43a9-8d52-ae44023847ba"));
+ private final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor();
+
+ private static final int RULE_DISPATCHER_POOL_SIZE = 2;
+ private static final int DB_CALLBACK_POOL_SIZE = 3;
+
@Mock
private TbContext ctxMock;
@Mock
@@ -95,8 +109,6 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
public void setUp() throws TbNodeException {
config = new CalculateDeltaNodeConfiguration().defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
- when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock);
-
node.init(ctxMock, nodeConfiguration);
}
@@ -110,6 +122,49 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
assertTrue(config.isTellFailureIfDeltaIsNegative());
}
+
+ @ParameterizedTest
+ @NullAndEmptySource
+ @ValueSource(strings = {" "}) // blank value
+ public void givenInvalidInputKey_whenInitThenThrowException(String key) {
+ config.setInputValueKey(key);
+ nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
+ var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration));
+ assertThat(exception).hasMessage("Input value key should be specified!");
+ assertThat(exception.isUnrecoverable()).isTrue();
+ }
+
+ @ParameterizedTest
+ @NullAndEmptySource
+ @ValueSource(strings = {" "}) // blank value
+ public void givenInvalidOutputKey_whenInitThenThrowException(String key) {
+ config.setOutputValueKey(key);
+ nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
+ var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration));
+ assertThat(exception).hasMessage("Output value key should be specified!");
+ assertThat(exception.isUnrecoverable()).isTrue();
+ }
+
+ @ParameterizedTest
+ @NullAndEmptySource
+ @ValueSource(strings = {" "}) // blank value
+ public void givenInvalidPeriodKey_whenInitThenThrowException(String key) {
+ config.setPeriodValueKey(key);
+ config.setAddPeriodBetweenMsgs(true);
+ nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
+ var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration));
+ assertThat(exception).hasMessage("Period value key should be specified!");
+ assertThat(exception.isUnrecoverable()).isTrue();
+ }
+
+ @Test
+ public void givenInvalidPeriodKeyAndAddPeriodDisabled_whenInitThenNoExceptionThrown() {
+ config.setPeriodValueKey(null);
+ config.setAddPeriodBetweenMsgs(false);
+ nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
+ assertDoesNotThrow(() -> node.init(ctxMock, nodeConfiguration));
+ }
+
@Test
public void givenInvalidMsgType_whenOnMsg_thenShouldTellNextOther() {
// GIVEN
@@ -120,7 +175,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
node.onMsg(ctxMock, msg);
// THEN
- verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
+ verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
verify(ctxMock, never()).tellSuccess(any());
verify(ctxMock, never()).tellFailure(any(), any());
}
@@ -134,7 +189,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
node.onMsg(ctxMock, msg);
// THEN
- verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
+ verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
verify(ctxMock, never()).tellSuccess(any());
verify(ctxMock, never()).tellFailure(any(), any());
}
@@ -149,7 +204,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
node.onMsg(ctxMock, msg);
// THEN
- verify(ctxMock, times(1)).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
+ verify(ctxMock).tellNext(eq(msg), eq(TbNodeConnectionType.OTHER));
verify(ctxMock, never()).tellSuccess(any());
verify(ctxMock, never()).tellFailure(any(), any());
}
@@ -175,7 +230,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -205,7 +260,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -235,7 +290,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -256,7 +311,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
node.init(ctxMock, nodeConfiguration);
- mockFindLatest(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", 40.0)));
+ mockFindLatestAsync(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", 40.0)));
var msgData = "{\"temperature\": 42,\"airPressure\":123}";
var firstMsgMetaData = new TbMsgMetaData();
@@ -269,7 +324,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -283,6 +338,8 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
reset(ctxMock);
reset(timeseriesServiceMock);
+ when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
+
var secondMsgMetaData = new TbMsgMetaData();
secondMsgMetaData.putValue("ts", String.valueOf(6L));
var secondMsg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, secondMsgMetaData, msgData);
@@ -294,7 +351,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(timeseriesServiceMock, never()).findLatest(any(), any(), anyList());
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -324,7 +381,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
verify(ctxMock, never()).tellFailure(any(), any());
@@ -341,7 +398,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
node.init(ctxMock, nodeConfiguration);
- mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L)));
+ mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L)));
var msgData = "{\"pulseCounter\":\"123\"}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
@@ -353,7 +410,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class);
- verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture());
+ verify(ctxMock).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture());
verify(ctxMock, never()).tellSuccess(any());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
@@ -373,7 +430,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
node.init(ctxMock, nodeConfiguration);
- mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L)));
+ mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L)));
var msgData = "{\"pulseCounter\":\"123\"}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
@@ -384,7 +441,7 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
// THEN
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class);
- verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture());
+ verify(ctxMock).tellSuccess(actualMsgCaptor.capture());
verify(ctxMock, never()).tellFailure(any(), any());
verify(ctxMock, never()).tellNext(any(), anyString());
verify(ctxMock, never()).tellNext(any(), anySet());
@@ -396,13 +453,23 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenInvalidStringValue_whenOnMsg_thenException() {
// GIVEN
- mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high")));
+ mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high")));
var msgData = "{\"pulseCounter\":\"123\"}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
- // WHEN-THEN
- Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
+ // WHEN
+ node.onMsg(ctxMock, msg);
+
+ // THEN
+ ArgumentCaptor throwableCaptor = ArgumentCaptor.forClass(Throwable.class);
+
+ verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture());
+ verify(ctxMock, never()).tellSuccess(any());
+ verify(ctxMock, never()).tellNext(any(), anyString());
+ verify(ctxMock, never()).tellNext(any(), anySet());
+
+ assertThat(throwableCaptor.getValue())
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Calculation failed. Unable to parse value [high] of telemetry [pulseCounter] to Double");
}
@@ -410,13 +477,23 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenBooleanValue_whenOnMsg_thenException() {
// GIVEN
- mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false)));
+ mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false)));
var msgData = "{\"pulseCounter\":true}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
- // WHEN-THEN
- Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
+ // WHEN
+ node.onMsg(ctxMock, msg);
+
+ // THEN
+ ArgumentCaptor throwableCaptor = ArgumentCaptor.forClass(Throwable.class);
+
+ verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture());
+ verify(ctxMock, never()).tellSuccess(any());
+ verify(ctxMock, never()).tellNext(any(), anyString());
+ verify(ctxMock, never()).tellNext(any(), anySet());
+
+ assertThat(throwableCaptor.getValue())
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Calculation failed. Boolean values are not supported!");
}
@@ -424,38 +501,104 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenJsonValue_whenOnMsg_thenException() {
// GIVEN
- mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new JsonDataEntry("pulseCounter", "{\"isActive\":false}")));
+ mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new JsonDataEntry("pulseCounter", "{\"isActive\":false}")));
var msgData = "{\"pulseCounter\":{\"isActive\":true}}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
- // WHEN-THEN
- Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg))
+ // WHEN
+ node.onMsg(ctxMock, msg);
+
+ // THEN
+ ArgumentCaptor throwableCaptor = ArgumentCaptor.forClass(Throwable.class);
+
+ verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture());
+ verify(ctxMock, never()).tellSuccess(any());
+ verify(ctxMock, never()).tellNext(any(), anyString());
+ verify(ctxMock, never()).tellNext(any(), anySet());
+
+ assertThat(throwableCaptor.getValue())
.isInstanceOf(IllegalArgumentException.class)
.hasMessage("Calculation failed. JSON values are not supported!");
}
+ @Test
+ public void givenConcurrentAccess_whenOnMsg_thenGetFromDBInvokedOnce() throws TbNodeException, InterruptedException {
+ DBCallbackExecutor dbCallbackExecutor = new DBCallbackExecutor();
+ dbCallbackExecutor.init();
+
+ RuleDispatcherExecutor ruleEngineDispatcherExecutor = new RuleDispatcherExecutor();
+ ruleEngineDispatcherExecutor.init();
+
+ assertThat(RULE_DISPATCHER_POOL_SIZE).as("dispatcher pool size have to be > 1").isGreaterThan(1);
+
+ final TbContext ctx = mock(TbContext.class);
+ final TimeseriesService timeseriesService = mock(TimeseriesService.class);
+
+ when(ctx.getTimeseriesService()).thenReturn(timeseriesService);
+ when(ctx.getDbCallbackExecutor()).thenReturn(dbCallbackExecutor);
+ when(timeseriesService.findLatest(any(), any(), anyString())).thenReturn(Futures.immediateFuture(Optional.empty()));
+
+ final CalculateDeltaNodeConfiguration config = new CalculateDeltaNodeConfiguration().defaultConfiguration();
+ final TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
+ final CalculateDeltaNode node = spy(CalculateDeltaNode.class);
+
+ node.init(ctx, nodeConfiguration);
+
+ List tbMsgList = IntStream.range(0, RULE_DISPATCHER_POOL_SIZE * 2).mapToObj(x -> {
+ var msgData = "{\"pulseCounter\":" + 2 + "}";
+ return TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
+ }).toList();
+
+ CountDownLatch processingLatch = new CountDownLatch(tbMsgList.size());
+
+ willAnswer(invocation -> {
+ processingLatch.countDown();
+ return invocation.callRealMethod();
+ }).given(node).processMsgAsync(any(), any());
+
+ tbMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg)));
+
+ assertThat(processingLatch.await(5, TimeUnit.SECONDS)).as("await on processingLatch").isTrue();
+
+ verify(timeseriesService).findLatest(any(), any(), anyString());
+ await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> verify(ctx, times(tbMsgList.size())).tellSuccess(any()));
+ }
+
+ private static class RuleDispatcherExecutor extends AbstractListeningExecutor {
+ @Override
+ protected int getThreadPollSize() {
+ return RULE_DISPATCHER_POOL_SIZE;
+ }
+ }
+
+ private static class DBCallbackExecutor extends AbstractListeningExecutor {
+ @Override
+ protected int getThreadPollSize() {
+ return DB_CALLBACK_POOL_SIZE;
+ }
+ }
+
@ParameterizedTest
@MethodSource("CalculateDeltaTestConfig")
public void givenCalculateDeltaConfig_whenOnMsg_thenVerify(CalculateDeltaTestConfig testConfig) throws TbNodeException {
// GIVEN
- config.setTellFailureIfDeltaIsNegative(testConfig.isTellFailureIfDeltaIsNegative());
- config.setExcludeZeroDeltas(testConfig.isExcludeZeroDeltas());
+ config.setTellFailureIfDeltaIsNegative(testConfig.tellFailureIfDeltaIsNegative());
+ config.setExcludeZeroDeltas(testConfig.excludeZeroDeltas());
config.setInputValueKey("temperature");
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
node.init(ctxMock, nodeConfiguration);
- mockFindLatest(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", testConfig.getPrevValue())));
+ mockFindLatestAsync(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", testConfig.prevValue())));
- var msgData = "{\"temperature\":" + testConfig.getCurrentValue() + ",\"airPressure\":123}";
+ var msgData = "{\"temperature\":" + testConfig.currentValue() + ",\"airPressure\":123}";
var msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DUMMY_DEVICE_ORIGINATOR, TbMsgMetaData.EMPTY, msgData);
// WHEN
-
node.onMsg(ctxMock, msg);
// THEN
- testConfig.getVerificationMethod().accept(ctxMock, msg);
+ testConfig.verificationMethod().accept(ctxMock, msg);
}
private static Stream CalculateDeltaTestConfig() {
@@ -510,47 +653,18 @@ public class CalculateDeltaNodeTest extends AbstractRuleNodeUpgradeTest {
);
}
- @Data
- @RequiredArgsConstructor
- private static class CalculateDeltaTestConfig {
- private final boolean tellFailureIfDeltaIsNegative;
- private final boolean excludeZeroDeltas;
- private final double prevValue;
- private final double currentValue;
- private final BiConsumer verificationMethod;
- }
-
- private void mockFindLatest(TsKvEntry tsKvEntry) {
- when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
- when(timeseriesServiceMock.findLatestSync(
- eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey())))
- )).thenReturn(List.of(tsKvEntry));
+ private record CalculateDeltaTestConfig(boolean tellFailureIfDeltaIsNegative, boolean excludeZeroDeltas,
+ double prevValue, double currentValue,
+ BiConsumer verificationMethod) {
}
private void mockFindLatestAsync(TsKvEntry tsKvEntry) {
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
+ when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock);
when(timeseriesServiceMock.findLatest(
- eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey())))
- )).thenReturn(Futures.immediateFuture(List.of(tsKvEntry)));
- }
-
- @RequiredArgsConstructor
- private static class ListMatcher implements ArgumentMatcher> {
-
- private final List expectedList;
-
- @Override
- public boolean matches(List actualList) {
- if (actualList == expectedList) {
- return true;
- }
- if (actualList.size() != expectedList.size()) {
- return false;
- }
- return actualList.containsAll(expectedList);
- }
-
+ eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), eq(tsKvEntry.getKey())
+ )).thenReturn(Futures.immediateFuture(Optional.of(tsKvEntry)));
}
private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
diff --git a/ui-ngx/package.json b/ui-ngx/package.json
index a5304207f5..b8a8e4c9a9 100644
--- a/ui-ngx/package.json
+++ b/ui-ngx/package.json
@@ -43,7 +43,6 @@
"@ngrx/store": "^15.4.0",
"@ngrx/store-devtools": "^15.4.0",
"@ngx-translate/core": "^14.0.0",
- "@ngx-translate/http-loader": "^7.0.0",
"@svgdotjs/svg.filter.js": "^3.0.8",
"@svgdotjs/svg.js": "^3.2.0",
"@tinymce/tinymce-angular": "^7.0.0",
diff --git a/ui-ngx/src/app/app.component.ts b/ui-ngx/src/app/app.component.ts
index 479e0465e5..cf53929e54 100644
--- a/ui-ngx/src/app/app.component.ts
+++ b/ui-ngx/src/app/app.component.ts
@@ -27,10 +27,13 @@ import { LocalStorageService } from '@core/local-storage/local-storage.service';
import { DomSanitizer } from '@angular/platform-browser';
import { MatIconRegistry } from '@angular/material/icon';
import { combineLatest } from 'rxjs';
-import { selectIsAuthenticated, selectIsUserLoaded } from '@core/auth/auth.selectors';
-import { distinctUntilChanged, filter, map, skip } from 'rxjs/operators';
+import { getCurrentAuthState, selectIsAuthenticated, selectIsUserLoaded } from '@core/auth/auth.selectors';
+import { distinctUntilChanged, filter, map, skip, tap } from 'rxjs/operators';
import { AuthService } from '@core/auth/auth.service';
import { svgIcons, svgIconsUrl } from '@shared/models/icon.models';
+import { isEqual } from '@core/utils';
+import { ActionSettingsChangeLanguage } from '@core/settings/settings.actions';
+import { SETTINGS_KEY } from '@core/settings/settings.effects';
@Component({
selector: 'tb-root',
@@ -92,8 +95,16 @@ export class AppComponent implements OnInit {
this.store.pipe(select(selectIsUserLoaded))]
).pipe(
map(results => ({isAuthenticated: results[0], isUserLoaded: results[1]})),
- distinctUntilChanged(),
- filter((data) => data.isUserLoaded ),
+ filter((data) => data.isUserLoaded),
+ distinctUntilChanged((a, b) => isEqual(a, b)),
+ tap((data) => {
+ let userLang = getCurrentAuthState(this.store).userDetails?.additionalInfo?.lang ?? null;
+ if (!userLang && !data.isAuthenticated) {
+ const settings = this.storageService.getItem(SETTINGS_KEY);
+ userLang = settings?.userLang ?? null;
+ }
+ this.notifyUserLang(userLang);
+ }),
skip(1),
).subscribe((data) => {
this.authService.gotoDefaultPlace(data.isAuthenticated);
@@ -111,4 +122,8 @@ export class AppComponent implements OnInit {
}
}
+ private notifyUserLang(userLang: string) {
+ this.store.dispatch(new ActionSettingsChangeLanguage({userLang}));
+ }
+
}
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 1848733006..1fcd18bd93 100644
--- a/ui-ngx/src/app/core/api/widget-api.models.ts
+++ b/ui-ngx/src/app/core/api/widget-api.models.ts
@@ -30,7 +30,12 @@ import {
import { TimeService } from '../services/time.service';
import { DeviceService } from '../http/device.service';
import { UtilsService } from '@core/services/utils.service';
-import { SubscriptionTimewindow, Timewindow, WidgetTimewindow } from '@shared/models/time/time.models';
+import {
+ ComparisonDuration,
+ SubscriptionTimewindow,
+ Timewindow,
+ WidgetTimewindow
+} from '@shared/models/time/time.models';
import { EntityType } from '@shared/models/entity-type.models';
import { HttpErrorResponse } from '@angular/common/http';
import { RafService } from '@core/services/raf.service';
@@ -265,7 +270,7 @@ export interface WidgetSubscriptionOptions {
onTimewindowChangeFunction?: (timewindow: Timewindow) => Timewindow;
legendConfig?: LegendConfig;
comparisonEnabled?: boolean;
- timeForComparison?: moment_.unitOfTime.DurationConstructor;
+ timeForComparison?: ComparisonDuration;
comparisonCustomIntervalValue?: number;
decimals?: number;
units?: string;
diff --git a/ui-ngx/src/app/core/api/widget-subscription.ts b/ui-ngx/src/app/core/api/widget-subscription.ts
index 933d836ffd..dc29f36ee2 100644
--- a/ui-ngx/src/app/core/api/widget-subscription.ts
+++ b/ui-ngx/src/app/core/api/widget-subscription.ts
@@ -23,13 +23,13 @@ import {
WidgetSubscriptionOptions
} from '@core/api/widget-api.models';
import {
- DataKey,
+ DataKey, DataKeySettingsWithComparison,
DataSet,
DataSetHolder,
Datasource,
DatasourceData,
datasourcesHasAggregation,
- DatasourceType,
+ DatasourceType, isDataKeySettingsWithComparison,
LegendConfig,
LegendData,
LegendKey,
@@ -513,7 +513,7 @@ export class WidgetSubscription implements IWidgetSubscription {
this.configuredDatasources.forEach((datasource, datasourceIndex) => {
const additionalDataKeys: DataKey[] = [];
datasource.dataKeys.forEach((dataKey, dataKeyIndex) => {
- if (dataKey.settings.comparisonSettings && dataKey.settings.comparisonSettings.showValuesForComparison) {
+ if (isDataKeySettingsWithComparison(dataKey.settings) && dataKey.settings.comparisonSettings.showValuesForComparison) {
const additionalDataKey = deepClone(dataKey);
additionalDataKey.isAdditional = true;
additionalDataKey.origDataKeyIndex = dataKeyIndex;
@@ -1468,11 +1468,12 @@ export class WidgetSubscription implements IWidgetSubscription {
if (datasource.isAdditional) {
const origDatasource = this.datasourcePages[datasource.origDatasourceIndex].data[dIndex];
datasource.dataKeys.forEach((dataKey) => {
- if (dataKey.settings.comparisonSettings.color) {
+ const settings: DataKeySettingsWithComparison = dataKey.settings;
+ if (settings.comparisonSettings.color) {
dataKey.color = dataKey.settings.comparisonSettings.color;
}
const origDataKey = origDatasource.dataKeys[dataKey.origDataKeyIndex];
- origDataKey.settings.comparisonSettings.color = dataKey.color;
+ (origDataKey.settings as DataKeySettingsWithComparison).comparisonSettings.color = dataKey.color;
});
}
});
@@ -1523,7 +1524,8 @@ export class WidgetSubscription implements IWidgetSubscription {
const formattedData = flatFormattedData(formattedDataArray);
datasource.dataKeys.forEach((dataKey) => {
- if (this.comparisonEnabled && dataKey.isAdditional && dataKey.settings.comparisonSettings.comparisonValuesLabel) {
+ if (this.comparisonEnabled && dataKey.isAdditional && isDataKeySettingsWithComparison(dataKey.settings) &&
+ dataKey.settings.comparisonSettings.comparisonValuesLabel) {
dataKey.label = createLabelFromPattern(dataKey.settings.comparisonSettings.comparisonValuesLabel, formattedData);
} else {
if (this.comparisonEnabled && dataKey.isAdditional) {
diff --git a/ui-ngx/src/app/core/auth/auth.service.ts b/ui-ngx/src/app/core/auth/auth.service.ts
index 0cc0ca4567..ea3dfb4143 100644
--- a/ui-ngx/src/app/core/auth/auth.service.ts
+++ b/ui-ngx/src/app/core/auth/auth.service.ts
@@ -22,7 +22,7 @@ import { Observable, of, ReplaySubject, throwError } from 'rxjs';
import { catchError, map, mergeMap, tap } from 'rxjs/operators';
import { LoginRequest, LoginResponse, PublicLoginRequest } from '@shared/models/login.models';
-import { ActivatedRoute, Router, UrlTree } from '@angular/router';
+import { Router, UrlTree } from '@angular/router';
import { defaultHttpOptions, defaultHttpOptionsFromConfig, RequestConfig } from '../http/http-utils';
import { UserService } from '../http/user.service';
import { Store } from '@ngrx/store';
@@ -35,7 +35,6 @@ import {
} from './auth.actions';
import { getCurrentAuthState, getCurrentAuthUser } from './auth.selectors';
import { Authority } from '@shared/models/authority.enum';
-import { ActionSettingsChangeLanguage } from '@app/core/settings/settings.actions';
import { AuthPayload, AuthState, SysParams, SysParamsState } from '@core/auth/auth.models';
import { TranslateService } from '@ngx-translate/core';
import { AuthUser } from '@shared/models/user.model';
@@ -59,7 +58,6 @@ export class AuthService {
private userService: UserService,
private timeService: TimeService,
private router: Router,
- private route: ActivatedRoute,
private zone: NgZone,
private utils: UtilsService,
private translate: TranslateService,
@@ -419,14 +417,7 @@ export class AuthService {
this.loadSystemParams().subscribe(
(sysParams) => {
authPayload = {...authPayload, ...sysParams};
- let userLang;
- if (authPayload.userDetails.additionalInfo && authPayload.userDetails.additionalInfo.lang) {
- userLang = authPayload.userDetails.additionalInfo.lang;
- } else {
- userLang = null;
- }
loadUserSubject.next(authPayload);
- this.notifyUserLang(userLang);
loadUserSubject.complete();
},
(err) => {
@@ -607,10 +598,6 @@ export class AuthService {
this.store.dispatch(new ActionAuthAuthenticated(authPayload));
}
- private notifyUserLang(userLang: string) {
- this.store.dispatch(new ActionSettingsChangeLanguage({userLang}));
- }
-
private updateAndValidateToken(token, prefix, notify) {
let valid = false;
const tokenData = this.jwtHelper.decodeToken(token);
diff --git a/ui-ngx/src/app/core/core.module.ts b/ui-ngx/src/app/core/core.module.ts
index 2f713e405f..a72eb1ada7 100644
--- a/ui-ngx/src/app/core/core.module.ts
+++ b/ui-ngx/src/app/core/core.module.ts
@@ -31,7 +31,6 @@ import {
TranslateModule,
TranslateParser
} from '@ngx-translate/core';
-import { TranslateHttpLoader } from '@ngx-translate/http-loader';
import { TbMissingTranslationHandler } from './translate/missing-translate-handler';
import { MatButtonModule } from '@angular/material/button';
import { MAT_DIALOG_DEFAULT_OPTIONS, MatDialogConfig, MatDialogModule } from '@angular/material/dialog';
@@ -41,9 +40,10 @@ import { TranslateDefaultCompiler } from '@core/translate/translate-default-comp
import { WINDOW_PROVIDERS } from '@core/services/window.service';
import { HotkeyModule } from 'angular2-hotkeys';
import { TranslateDefaultParser } from '@core/translate/translate-default-parser';
+import { TranslateDefaultLoader } from '@core/translate/translate-default-loader';
export function HttpLoaderFactory(http: HttpClient) {
- return new TranslateHttpLoader(http, './assets/locale/locale.constant-', '.json');
+ return new TranslateDefaultLoader(http);
}
@NgModule({
diff --git a/ui-ngx/src/app/core/http/device-profile.service.ts b/ui-ngx/src/app/core/http/device-profile.service.ts
index 91ca903c7a..7a9dfd13ba 100644
--- a/ui-ngx/src/app/core/http/device-profile.service.ts
+++ b/ui-ngx/src/app/core/http/device-profile.service.ts
@@ -60,7 +60,7 @@ export class DeviceProfileService {
public getLwm2mObjects(sortOrder: SortOrder, objectIds?: string[], searchText?: string, config?: RequestConfig):
Observable> {
- let url = `/api/resource/lwm2m/?sortProperty=${sortOrder.property}&sortOrder=${sortOrder.direction}`;
+ let url = `/api/resource/lwm2m?sortProperty=${sortOrder.property}&sortOrder=${sortOrder.direction}`;
if (isDefinedAndNotNull(objectIds) && objectIds.length > 0) {
url += `&objectIds=${objectIds}`;
}
diff --git a/ui-ngx/src/app/core/services/resources.service.ts b/ui-ngx/src/app/core/services/resources.service.ts
index bd2402b44c..449e950a65 100644
--- a/ui-ngx/src/app/core/services/resources.service.ts
+++ b/ui-ngx/src/app/core/services/resources.service.ts
@@ -34,6 +34,7 @@ import { select, Store } from '@ngrx/store';
import { selectIsAuthenticated } from '@core/auth/auth.selectors';
import { AppState } from '@core/core.state';
import { map, tap } from 'rxjs/operators';
+import { RequestConfig } from '@core/http/http-utils';
declare const System;
@@ -106,11 +107,11 @@ export class ResourcesService {
return this.loadResourceByType(fileType, url);
}
- public downloadResource(downloadUrl: string): Observable {
- return this.http.get(downloadUrl, {
+ public downloadResource(downloadUrl: string, config?: RequestConfig): Observable {
+ return this.http.get(downloadUrl, {...config, ...{
responseType: 'arraybuffer',
observe: 'response'
- }).pipe(
+ }}).pipe(
map((response) => {
const headers = response.headers;
const filename = headers.get('x-filename');
diff --git a/ui-ngx/src/app/core/settings/settings.effects.ts b/ui-ngx/src/app/core/settings/settings.effects.ts
index 56615bf62c..1564814a74 100644
--- a/ui-ngx/src/app/core/settings/settings.effects.ts
+++ b/ui-ngx/src/app/core/settings/settings.effects.ts
@@ -28,7 +28,6 @@ import { AppState } from '@app/core/core.state';
import { LocalStorageService } from '@app/core/local-storage/local-storage.service';
import { TitleService } from '@app/core/services/title.service';
import { updateUserLang } from '@app/core/settings/settings.utils';
-import { AuthService } from '@core/auth/auth.service';
import { UtilsService } from '@core/services/utils.service';
import { getCurrentAuthUser } from '@core/auth/auth.selectors';
import { ActionAuthUpdateLastPublicDashboardId } from '../auth/auth.actions';
@@ -40,7 +39,6 @@ export class SettingsEffects {
constructor(
private actions$: Actions,
private store: Store,
- private authService: AuthService,
private utils: UtilsService,
private router: Router,
private localStorageService: LocalStorageService,
@@ -49,26 +47,19 @@ export class SettingsEffects {
) {
}
-
- persistSettings = createEffect(() => this.actions$.pipe(
+ setTranslateServiceLanguage = createEffect(() => this.actions$.pipe(
ofType(
SettingsActionTypes.CHANGE_LANGUAGE,
),
withLatestFrom(this.store.pipe(select(selectSettingsState))),
- tap(([action, settings]) =>
- this.localStorageService.setItem(SETTINGS_KEY, settings)
- )
- ), {dispatch: false});
-
-
- setTranslateServiceLanguage = createEffect(() => this.store.pipe(
- select(selectSettingsState),
- map(settings => settings.userLang),
- distinctUntilChanged(),
- tap(userLang => updateUserLang(this.translate, userLang))
+ map(settings => settings[1]),
+ distinctUntilChanged((a, b) => a?.userLang === b?.userLang),
+ tap(setting => {
+ this.localStorageService.setItem(SETTINGS_KEY, setting);
+ updateUserLang(this.translate, setting.userLang);
+ })
), {dispatch: false});
-
setTitle = createEffect(() => merge(
this.actions$.pipe(ofType(SettingsActionTypes.CHANGE_LANGUAGE)),
this.router.events.pipe(filter(event => event instanceof ActivationEnd))
@@ -81,7 +72,6 @@ export class SettingsEffects {
})
), {dispatch: false});
-
setPublicId = createEffect(() => merge(
this.router.events.pipe(filter(event => event instanceof ActivationEnd))
).pipe(
diff --git a/ui-ngx/src/app/core/settings/settings.utils.ts b/ui-ngx/src/app/core/settings/settings.utils.ts
index dfe8b37519..e816bbffb0 100644
--- a/ui-ngx/src/app/core/settings/settings.utils.ts
+++ b/ui-ngx/src/app/core/settings/settings.utils.ts
@@ -18,7 +18,7 @@ import { environment as env } from '@env/environment';
import { TranslateService } from '@ngx-translate/core';
import * as _moment from 'moment';
-export function updateUserLang(translate: TranslateService, userLang: string) {
+export function updateUserLang(translate: TranslateService, userLang: string, translations = env.supportedLangs) {
let targetLang = userLang;
if (!env.production) {
console.log(`User lang: ${targetLang}`);
@@ -29,7 +29,7 @@ export function updateUserLang(translate: TranslateService, userLang: string) {
console.log(`Fallback to browser lang: ${targetLang}`);
}
}
- const detectedSupportedLang = detectSupportedLang(targetLang);
+ const detectedSupportedLang = detectSupportedLang(targetLang, translations);
if (!env.production) {
console.log(`Detected supported lang: ${detectedSupportedLang}`);
}
@@ -37,10 +37,10 @@ export function updateUserLang(translate: TranslateService, userLang: string) {
_moment.locale([detectedSupportedLang]);
}
-function detectSupportedLang(targetLang: string): string {
+function detectSupportedLang(targetLang: string, translations: string[]): string {
const langTag = (targetLang || '').split('-').join('_');
if (langTag.length) {
- if (env.supportedLangs.indexOf(langTag) > -1) {
+ if (translations.indexOf(langTag) > -1) {
return langTag;
} else {
const parts = langTag.split('_');
@@ -50,7 +50,7 @@ function detectSupportedLang(targetLang: string): string {
} else {
lang = langTag;
}
- const foundLangs = env.supportedLangs.filter(
+ const foundLangs = translations.filter(
(supportedLang: string) => {
const supportedLangParts = supportedLang.split('_');
return supportedLangParts[0] === lang;
diff --git a/ui-ngx/src/app/core/translate/translate-default-loader.ts b/ui-ngx/src/app/core/translate/translate-default-loader.ts
new file mode 100644
index 0000000000..7f26170c9c
--- /dev/null
+++ b/ui-ngx/src/app/core/translate/translate-default-loader.ts
@@ -0,0 +1,30 @@
+///
+/// Copyright © 2016-2024 The Thingsboard Authors
+///
+/// Licensed under the Apache License, Version 2.0 (the "License");
+/// you may not use this file except in compliance with the License.
+/// You may obtain a copy of the License at
+///
+/// http://www.apache.org/licenses/LICENSE-2.0
+///
+/// Unless required by applicable law or agreed to in writing, software
+/// distributed under the License is distributed on an "AS IS" BASIS,
+/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+/// See the License for the specific language governing permissions and
+/// limitations under the License.
+///
+
+import { TranslateLoader } from '@ngx-translate/core';
+import { Observable } from 'rxjs';
+import { HttpClient } from '@angular/common/http';
+
+export class TranslateDefaultLoader implements TranslateLoader {
+
+ constructor(private http: HttpClient) {
+
+ }
+
+ getTranslation(lang: string): Observable