Browse Source

Merge branch 'rc' into features/add_tooltip_option_to_show_stack_mode_total_value_on_timeseries_chart_widgets

pull/13373/head
Paolo Cristiani 1 year ago
committed by GitHub
parent
commit
938c768fb3
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 13
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java
  3. 8
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  4. 11
      dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java
  5. 19
      dao/src/main/java/org/thingsboard/server/dao/util/TenantRateLimitException.java
  6. 16
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java
  7. 19
      netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java
  8. 20
      netty-mqtt/src/test/java/org/thingsboard/mqtt/MqttClientTest.java

13
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java

@ -16,13 +16,16 @@
package org.thingsboard.server.service.cf.ctx.state;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.core.type.TypeReference;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.script.api.tbel.TbelCfArg;
import org.thingsboard.script.api.tbel.TbelCfSingleValueArg;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BasicKvEntry;
import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.util.ProtoUtils;
@ -90,7 +93,15 @@ public class SingleValueArgumentEntry implements ArgumentEntry {
@Override
public TbelCfArg toTbelCfArg() {
return new TbelCfSingleValueArg(ts, kvEntryValue.getValue());
Object value = kvEntryValue.getValue();
if (kvEntryValue instanceof JsonDataEntry) {
try {
value = JacksonUtil.readValue(kvEntryValue.getValueAsString(), new TypeReference<>() {
});
} catch (Exception e) {
}
}
return new TbelCfSingleValueArg(ts, value);
}
@Override

4
application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java

@ -36,6 +36,7 @@ import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.exception.RateLimitExceededException;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -52,7 +53,6 @@ import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.util.TenantRateLimitException;
import org.thingsboard.server.exception.UnauthorizedException;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.util.TbCoreComponent;
@ -742,7 +742,7 @@ public class DefaultWebSocketService implements WebSocketService {
@Override
public void onFailure(Throwable e) {
if (e instanceof TenantRateLimitException || e.getCause() instanceof TenantRateLimitException) {
if (e instanceof RateLimitExceededException || e.getCause() instanceof RateLimitExceededException) {
log.trace("[{}] Tenant rate limit detected for subscription: [{}]:{}", sessionRef.getSecurityCtx().getTenantId(), entityId, cmd);
} else {
log.info(FAILED_TO_FETCH_DATA, e);

8
dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java

@ -205,8 +205,8 @@ public class BaseTimeseriesService implements TimeseriesService {
ListenableFuture<Integer> dpsFuture = saveTs ? Futures.transform(Futures.allAsList(tsFutures), SUM_ALL_INTEGERS, MoreExecutors.directExecutor()) : Futures.immediateFuture(0);
ListenableFuture<List<Long>> versionsFuture = saveLatest ? Futures.allAsList(latestFutures) : Futures.immediateFuture(null);
return Futures.whenAllComplete(dpsFuture, versionsFuture).call(() -> {
Integer dataPoints = Futures.getUnchecked(dpsFuture);
List<Long> versions = Futures.getUnchecked(versionsFuture);
Integer dataPoints = dpsFuture.get();
List<Long> versions = versionsFuture.get();
return TimeseriesSaveResult.of(dataPoints, versions);
}, MoreExecutors.directExecutor());
}
@ -298,13 +298,13 @@ public class BaseTimeseriesService implements TimeseriesService {
long interval = query.getInterval();
if (interval < 1) {
throw new IncorrectParameterException("Invalid TsKvQuery: 'interval' must be greater than 0, but got " + interval +
". Please check your query parameters and ensure 'endTs' is greater than 'startTs' or increase 'interval'.");
". Please check your query parameters and ensure 'endTs' is greater than 'startTs' or increase 'interval'.");
}
long step = Math.max(interval, 1000);
long intervalCounts = (query.getEndTs() - query.getStartTs()) / step;
if (intervalCounts > maxTsIntervals || intervalCounts < 0) {
throw new IncorrectParameterException("Incorrect TsKvQuery. Number of intervals is to high - " + intervalCounts + ". " +
"Please increase 'interval' parameter for your query or reduce the time range of the query.");
"Please increase 'interval' parameter for your query or reduce the time range of the query.");
}
}
}

11
dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java

@ -32,6 +32,7 @@ import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.cache.limits.RateLimitService;
import org.thingsboard.server.common.data.exception.RateLimitExceededException;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.msg.queue.ServiceType;
@ -66,7 +67,7 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
private final long maxWaitTime;
private final long pollMs;
private final String bufferName;
private final String bufferName;
private final BlockingQueue<AsyncTaskContext<T, V>> queue;
private final ExecutorService dispatcherExecutor;
private final ExecutorService callbackExecutor;
@ -124,7 +125,7 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
if (!rateLimitService.checkRateLimit(myLimitedApi, tenantId, tenantId, true)) {
stats.incrementRateLimitedTenant(tenantId);
stats.getTotalRateLimited().increment();
settableFuture.setException(new TenantRateLimitException());
settableFuture.setException(new RateLimitExceededException(myLimitedApi));
perTenantLimitReached = true;
}
} else if (tenantId == null) {
@ -299,9 +300,9 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
.count();
if (queueSize > 0
|| rateLimitedTenantsCount > 0
|| concurrencyLevel.get() > 0
|| stats.getStatsCounters().stream().anyMatch(counter -> counter.get() > 0)
|| rateLimitedTenantsCount > 0
|| concurrencyLevel.get() > 0
|| stats.getStatsCounters().stream().anyMatch(counter -> counter.get() > 0)
) {
StringBuilder statsBuilder = new StringBuilder();

19
dao/src/main/java/org/thingsboard/server/dao/util/TenantRateLimitException.java

@ -1,19 +0,0 @@
/**
* Copyright © 2016-2025 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.dao.util;
public class TenantRateLimitException extends Exception {
}

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

@ -556,22 +556,6 @@ public class MqttClientTest extends AbstractContainerTest {
assertThat(provisionResponse.get("status").asText()).isEqualTo("NOT_FOUND");
}
@Test
public void regularDisconnect() throws Exception {
DeviceCredentials deviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(device.getId());
MqttMessageListener listener = new MqttMessageListener();
MqttClient mqttClient = getMqttClient(deviceCredentials, listener, MqttVersion.MQTT_5);
final List<Byte> returnCodeByteValue = new ArrayList<>();
MqttClientCallback callbackForDisconnectWithReturnCode = getCallbackWrapperForDisconnectWithReturnCode(returnCodeByteValue);
mqttClient.setCallback(callbackForDisconnectWithReturnCode);
mqttClient.disconnect();
Thread.sleep(1000);
assertThat(returnCodeByteValue.size()).isEqualTo(1);
MqttReasonCodes.Disconnect returnCode = MqttReasonCodes.Disconnect.valueOf(returnCodeByteValue.get(0));
assertThat(returnCode).isEqualTo(MqttReasonCodes.Disconnect.NORMAL_DISCONNECT);
}
@Test
public void clientSessionTakenOverDisconnect() throws Exception {
DeviceCredentials deviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(device.getId());

19
netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java

@ -96,6 +96,8 @@ final class MqttClientImpl implements MqttClient {
private final ListeningExecutor handlerExecutor;
private final static int DISCONNECT_FALLBACK_DELAY_SECS = 1;
/**
* Construct the MqttClientImpl with default config
*/
@ -456,16 +458,25 @@ final class MqttClientImpl implements MqttClient {
@Override
public void disconnect() {
if (disconnected) {
return;
}
log.trace("[{}] Disconnecting from server", channel != null ? channel.id() : "UNKNOWN");
disconnected = true;
if (this.channel != null) {
MqttMessage message = new MqttMessage(new MqttFixedHeader(MqttMessageType.DISCONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0));
ChannelFuture channelFuture = this.sendAndFlushPacket(message);
sendAndFlushPacket(message).addListener((ChannelFutureListener) future -> {
future.channel().close();
disconnected = true;
});
eventLoop.schedule(() -> {
if (!channelFuture.isDone()) {
if (channel.isOpen()) {
log.trace("[{}] Channel still open after {} second; forcing close now", channel.id(), DISCONNECT_FALLBACK_DELAY_SECS);
this.channel.close();
disconnected = true;
}
}, 500, TimeUnit.MILLISECONDS);
}, DISCONNECT_FALLBACK_DELAY_SECS, TimeUnit.SECONDS);
}
}

20
netty-mqtt/src/test/java/org/thingsboard/mqtt/MqttClientTest.java

@ -119,6 +119,26 @@ class MqttClientTest {
assertThat(client.isConnected()).isTrue();
}
@Test
void testDisconnectFromBroker() {
// GIVEN
var clientConfig = new MqttClientConfig();
clientConfig.setOwnerId("Test[Disconnect]");
clientConfig.setClientId("disconnect");
client = MqttClient.create(clientConfig, null, handlerExecutor);
connect(broker.getHost(), broker.getMqttPort());
// WHEN
client.disconnect();
// THEN
Awaitility.await("waiting for client to disconnect")
.atMost(Duration.ofSeconds(5))
.untilAsserted(() -> assertThat(client.isConnected()).isFalse());
}
@Test
void testDisconnectDueToKeepAliveIfNoActivity() {
// GIVEN

Loading…
Cancel
Save