Browse Source

Merge remote-tracking branch 'origin/lts-4.2' into lts-4.3

pull/15439/head
Viacheslav Klimov 7 months ago
parent
commit
6ce51c6277
Failed to extract signature
  1. 40
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  2. 23
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  3. 11
      application/src/main/resources/thingsboard.yml
  4. 12
      application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java
  5. 36
      application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java
  6. 272
      application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java
  7. 8
      application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java
  8. 4
      common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java
  9. 1
      common/edge-api/src/main/proto/edge.proto
  10. 2
      dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java

40
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -19,7 +19,10 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import io.grpc.Server;
import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts;
import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder;
import io.grpc.netty.shaded.io.netty.handler.ssl.SslContext;
import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder;
import io.grpc.stub.StreamObserver;
import jakarta.annotation.Nullable;
import jakarta.annotation.PreDestroy;
@ -37,7 +40,8 @@ import org.thingsboard.server.cache.TbTransactionalCache;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.ResourceUtils;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.transport.config.ssl.PemSslCredentials;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.id.EdgeId;
@ -67,7 +71,6 @@ import org.thingsboard.server.service.edge.EdgeContextComponent;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import java.io.IOException;
import java.io.InputStream;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
@ -112,6 +115,8 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private String certFileResource;
@Value("${edges.rpc.ssl.private_key}")
private String privateKeyResource;
@Value("${edges.rpc.ssl.key_password:}")
private String keyPassword;
@Value("${edges.state.persistToTelemetry:false}")
private boolean persistToTelemetry;
@Value("${edges.rpc.client_max_keep_alive_time_sec:1}")
@ -176,9 +181,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
.addService(this);
if (sslEnabled) {
try {
InputStream certFileIs = ResourceUtils.getInputStream(this, certFileResource);
InputStream privateKeyFileIs = ResourceUtils.getInputStream(this, privateKeyResource);
builder.useTransportSecurity(certFileIs, privateKeyFileIs);
setupSsl(builder);
} catch (Exception e) {
log.error("Unable to set up SSL context. Reason: " + e.getMessage(), e);
throw new RuntimeException("Unable to set up SSL context!", e);
@ -199,6 +202,33 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.info("Edge RPC service initialized!");
}
/**
* Configures TLS for the Edge gRPC server.
* <p>
* Delegates PEM parsing and key management to {@link PemSslCredentials} — the same
* class used by MQTT, CoAP, and LwM2M transports — which supports:
* <ul>
* <li>Separate certificate and private key files (classic two-file setup)</li>
* <li>Combined PEM: certificate chain + private key in a single {@code cert} file
* ({@code private_key} left empty)</li>
* <li>Encrypted private keys (password supplied via {@code key_password})</li>
* </ul>
* Path resolution (for both {@code cert} and {@code private_key}) is handled by
* {@link org.thingsboard.server.common.data.ResourceUtils#getInputStream ResourceUtils}:
* absolute path → relative / working-dir → classpath → {@code classpath:} prefix.
*/
void setupSsl(NettyServerBuilder builder) throws Exception {
PemSslCredentials credentials = new PemSslCredentials();
credentials.setCertFile(certFileResource);
credentials.setKeyFile(StringUtils.isEmpty(privateKeyResource) ? null : privateKeyResource);
credentials.setKeyPassword(keyPassword);
credentials.init(false);
SslContext sslContext = GrpcSslContexts.configure(
SslContextBuilder.forServer(credentials.createKeyManagerFactory())).build();
builder.sslContext(sslContext);
}
@PreDestroy
public void destroy() {
if (server != null) {

23
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -34,6 +34,7 @@ import org.springframework.web.socket.CloseStatus;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
@ -242,7 +243,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
@Override
public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId());
log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId(), t);
if (t instanceof ResultSetSizeLimitExceededException) {
sendError(finalCtx, t);
}
}
}, wsCallBackExecutor);
}
@ -258,7 +262,18 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
handleLatestCmd(ctx, cmd.getLatestCmd());
}
if (cmd.getTsCmd() != null) {
handleTimeSeriesCmd(ctx, cmd.getTsCmd());
Futures.addCallback(handleTimeSeriesCmd(ctx, cmd.getTsCmd()), new FutureCallback<>() {
@Override
public void onSuccess(TbEntityDataSubCtx result) {}
@Override
public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to process timeseries command", ctx.getSessionId(), ctx.getCmdId(), t);
if (t instanceof ResultSetSizeLimitExceededException) {
sendError(ctx, t);
}
}
}, wsCallBackExecutor);
}
} else {
checkAndSendInitialData(ctx);
@ -268,6 +283,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
}
}
private void sendError(TbEntityDataSubCtx ctx, Throwable t) {
ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage()));
}
private void checkAndSendInitialData(@Nullable TbEntityDataSubCtx theCtx) {
if (!theCtx.isInitialDataSent()) {
EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription());

11
application/src/main/resources/thingsboard.yml

@ -1505,10 +1505,17 @@ edges:
ssl:
# Enable/disable SSL support
enabled: "${EDGES_RPC_SSL_ENABLED:false}"
# Cert file to be used during TLS connectivity to the cloud
# Path to the server certificate file (holds server certificate or certificate chain, may include server private key).
# Accepts an absolute filesystem path (e.g. /etc/thingsboard/certChainFile.pem),
# a relative path resolved against the working directory first then the classpath,
# or a classpath resource with the explicit "classpath:" prefix (e.g. classpath:conf/certChainFile.pem).
cert: "${EDGES_RPC_SSL_CERT:certChainFile.pem}"
# Private key file associated with the Cert certificate. This key is used in the encryption process during a secure connection
# Path to the server certificate private key file. Optional if the private key is already present in the cert file above.
# Supports the same path resolution as 'cert': absolute, relative/classpath, or "classpath:" prefix.
# Leave empty when using a combined PEM cert that already contains the private key.
private_key: "${EDGES_RPC_SSL_PRIVATE_KEY:privateKeyFile.pem}"
# Server certificate private key password (optional). Leave empty if the key is not encrypted.
key_password: "${EDGES_RPC_SSL_KEY_PASSWORD:}"
# Maximum size (in bytes) of inbound messages the cloud can handle from the edge. By default, it can handle messages up to 4 Megabytes
max_inbound_message_size: "${EDGES_RPC_MAX_INBOUND_MESSAGE_SIZE:4194304}"
# Maximum length of telemetry (time-series and attributes) message the cloud sends to the edge. By default, there is no limitation.

12
application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java

@ -41,6 +41,7 @@ import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.util.TestSocketUtils;
import org.springframework.test.web.servlet.ResultActions;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.Customer;
@ -90,8 +91,6 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.MQTT_PORT;
import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.MQTT_URL;
@TestPropertySource(properties = {
"transport.mqtt.enabled=true",
@ -101,6 +100,15 @@ import static org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest.
@ContextConfiguration(classes = {EntityViewControllerTest.Config.class})
@DaoSqlTest
public class EntityViewControllerTest extends AbstractControllerTest {
// Must NOT be imported from AbstractMqttIntegrationTest. That field is a static final initialized
// once per JVM. Other test classes (e.g. MqttGatewayRateLimitsTest, DeviceEdgeTest) share the same
// constant but produce a different Spring context cache key, so Spring creates a separate
// ApplicationContext for each of them. Every context starts its own MqttTransportService and tries
// to bind the same port, causing BindException when tests run in the same Surefire JVM fork.
// Declaring the port here gives this context its own independently allocated port.
static final int MQTT_PORT = TestSocketUtils.findAvailableTcpPort();
static final String MQTT_URL = "tcp://localhost:" + MQTT_PORT;
@DynamicPropertySource
static void props(DynamicPropertyRegistry registry) {
log.warn("transport.mqtt.bind_port = {}", MQTT_PORT);

36
application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java

@ -19,13 +19,16 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource;
import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils;
import org.thingsboard.common.util.JacksonUtil;
@ -60,7 +63,9 @@ import org.thingsboard.server.common.data.query.NumericFilterPredicate;
import org.thingsboard.server.common.data.query.SingleEntityFilter;
import org.thingsboard.server.common.data.query.TsValue;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.service.subscription.SubscriptionErrorCode;
import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
@ -95,6 +100,9 @@ public class WebsocketApiTest extends AbstractControllerTest {
@Autowired
private TelemetrySubscriptionService tsService;
@SpyBean
private TimeseriesService timeseriesService;
Device device;
DeviceTypeFilter dtf;
@ -965,6 +973,34 @@ public class WebsocketApiTest extends AbstractControllerTest {
}
@Test
public void testHistoryCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception {
ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L);
Mockito.doReturn(Futures.immediateFailedFuture(exception))
.when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any());
List<String> keys = List.of("temperature");
long now = System.currentTimeMillis();
EntityDataUpdate errorUpdate = getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf);
assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode());
assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage());
}
@Test
public void testTimeSeriesCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception {
ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L);
Mockito.doReturn(Futures.immediateFailedFuture(exception))
.when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any());
List<String> keys = List.of("temperature");
long now = System.currentTimeMillis();
EntityDataUpdate errorUpdate = getWsClient().subscribeTsUpdate(keys, now, TimeUnit.HOURS.toMillis(1), dtf);
assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode());
assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage());
}
private void sendTelemetry(Device device, List<TsKvEntry> tsData) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(1);
tsService.saveTimeseries(TimeseriesSaveRequest.builder()

272
application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java

@ -0,0 +1,272 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.edge.rpc;
import io.grpc.ManagedChannel;
import io.grpc.Server;
import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts;
import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder;
import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder;
import org.bouncycastle.asn1.x500.X500Name;
import org.bouncycastle.cert.jcajce.JcaX509CertificateConverter;
import org.bouncycastle.cert.jcajce.JcaX509v3CertificateBuilder;
import org.bouncycastle.jce.provider.BouncyCastleProvider;
import org.bouncycastle.openssl.jcajce.JcaPEMWriter;
import org.bouncycastle.openssl.jcajce.JcePEMEncryptorBuilder;
import org.bouncycastle.operator.jcajce.JcaContentSignerBuilder;
import org.bouncycastle.util.io.pem.PemObject;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.EnumSource;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.controller.AbstractWebTest;
import org.thingsboard.server.gen.edge.v1.EdgeRpcServiceGrpc;
import java.io.ByteArrayInputStream;
import java.math.BigInteger;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.KeyPair;
import java.security.KeyPairGenerator;
import java.security.PrivateKey;
import java.security.Security;
import java.security.cert.X509Certificate;
import java.security.spec.ECGenParameterSpec;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.awaitility.Awaitility.await;
/**
* Tests for Edge gRPC SSL setup using the production {@link EdgeGrpcService#setupSsl} method.
* <p>
* Covers:
* 1. Separate cert and key PEM inputs
* 2. Combined PEM (cert + key in one file)
* 3. Encrypted private key with password
* 4. Missing key in combined PEM → error
* <p>
* Each scenario is parameterized across key types: RSA-2048, RSA-4096, EC P-256, EC P-384.
*/
class EdgeGrpcSslTest {
static {
if (Security.getProvider(BouncyCastleProvider.PROVIDER_NAME) == null) {
Security.addProvider(new BouncyCastleProvider());
}
}
enum KeyType {
RSA_2048("RSA", 2048, null, "SHA256withRSA"),
RSA_4096("RSA", 4096, null, "SHA256withRSA"),
EC_P256("EC", 256, "secp256r1", "SHA256withECDSA"),
EC_P384("EC", 384, "secp384r1", "SHA384withECDSA");
final String algorithm;
final int size;
final String curve;
final String sigAlg;
KeyType(String algorithm, int size, String curve, String sigAlg) {
this.algorithm = algorithm;
this.size = size;
this.curve = curve;
this.sigAlg = sigAlg;
}
KeyPair generateKeyPair() throws Exception {
KeyPairGenerator kpg = KeyPairGenerator.getInstance(algorithm);
if (curve != null) {
kpg.initialize(new ECGenParameterSpec(curve));
} else {
kpg.initialize(size);
}
return kpg.generateKeyPair();
}
}
private final List<Path> tempFiles = new ArrayList<>();
private Server server;
private ManagedChannel channel;
@AfterEach
void cleanup() throws Exception {
if (channel != null) {
channel.shutdownNow().awaitTermination(2, TimeUnit.SECONDS);
}
if (server != null) {
server.shutdownNow().awaitTermination(2, TimeUnit.SECONDS);
}
for (Path p : tempFiles) {
Files.deleteIfExists(p);
}
}
@ParameterizedTest(name = "separateCertAndKey_{0}")
@EnumSource(KeyType.class)
void separateCertAndKey(KeyType keyType) throws Exception {
KeyPair kp = keyType.generateKeyPair();
X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
Path certFile = writeTempPem("cert", cert);
Path keyFile = writeTempPem("key", kp.getPrivate());
server = startServer(certFile.toString(), keyFile.toString(), null);
assertTlsConnectivity(cert);
}
@ParameterizedTest(name = "combinedPemWithCertAndKey_{0}")
@EnumSource(KeyType.class)
void combinedPemWithCertAndKey(KeyType keyType) throws Exception {
KeyPair kp = keyType.generateKeyPair();
X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
Path combinedFile = writeTempPem("combined", cert, kp.getPrivate());
server = startServer(combinedFile.toString(), "", null);
assertTlsConnectivity(cert);
}
// RSA-only: BouncyCastle writes encrypted EC keys in traditional PEM format (BEGIN EC PRIVATE KEY),
// which after decryption produces a PEMKeyPair without public key info — causing PemSslCredentials
// to fail with "Cannot invoke SubjectPublicKeyInfo.getEncoded() because getPublicKeyInfo() is null".
@ParameterizedTest(name = "encryptedPrivateKey_{0}")
@EnumSource(value = KeyType.class, names = {"RSA_2048", "RSA_4096"})
void encryptedPrivateKey(KeyType keyType) throws Exception {
KeyPair kp = keyType.generateKeyPair();
X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
String password = "test-password";
Path combinedFile = writeTempPemEncrypted("enc-combined", password, cert, kp.getPrivate());
server = startServer(combinedFile.toString(), "", password);
assertTlsConnectivity(cert);
}
@ParameterizedTest(name = "combinedPemWithCertOnly_throwsException_{0}")
@EnumSource(KeyType.class)
void combinedPemWithCertOnly_throwsException(KeyType keyType) throws Exception {
KeyPair kp = keyType.generateKeyPair();
X509Certificate cert = generateSelfSignedCert(kp, keyType.sigAlg);
Path certOnlyFile = writeTempPem("cert-only", cert);
assertThatThrownBy(() -> startServer(certOnlyFile.toString(), "", null))
.isInstanceOf(IllegalArgumentException.class);
}
// --- Server startup using production EdgeGrpcService.setupSsl() ---
private Server startServer(String certFileResource, String privateKeyResource, String keyPassword) throws Exception {
EdgeGrpcService edgeGrpcService = new EdgeGrpcService();
ReflectionTestUtils.setField(edgeGrpcService, "certFileResource", certFileResource);
ReflectionTestUtils.setField(edgeGrpcService, "privateKeyResource", privateKeyResource);
ReflectionTestUtils.setField(edgeGrpcService, "keyPassword", keyPassword != null ? keyPassword : "");
NettyServerBuilder builder = NettyServerBuilder.forPort(0)
.addService(new EdgeRpcServiceGrpc.EdgeRpcServiceImplBase() {});
edgeGrpcService.setupSsl(builder);
return builder.build().start();
}
private void assertTlsConnectivity(X509Certificate trustedCert) throws Exception {
String certPem = toPem(trustedCert);
var clientSsl = GrpcSslContexts.forClient()
.trustManager(new ByteArrayInputStream(certPem.getBytes(StandardCharsets.UTF_8)))
.build();
channel = NettyChannelBuilder.forAddress("localhost", server.getPort())
.sslContext(clientSsl)
.build();
channel.getState(true); // trigger connection attempt
await().atMost(AbstractWebTest.TIMEOUT, TimeUnit.SECONDS)
.pollInterval(50, TimeUnit.MILLISECONDS)
.untilAsserted(() -> {
var state = channel.getState(false);
if (state == io.grpc.ConnectivityState.TRANSIENT_FAILURE) {
throw new AssertionError("TLS handshake failed: channel in TRANSIENT_FAILURE");
}
assertThat(state).isEqualTo(io.grpc.ConnectivityState.READY);
});
}
// --- Cert/key generation ---
private X509Certificate generateSelfSignedCert(KeyPair kp, String sigAlg) throws Exception {
X500Name subject = new X500Name("CN=localhost");
Date now = new Date();
return new JcaX509CertificateConverter().getCertificate(
new JcaX509v3CertificateBuilder(
subject, BigInteger.ONE, now,
new Date(now.getTime() + TimeUnit.DAYS.toMillis(1)),
subject, kp.getPublic())
.build(new JcaContentSignerBuilder(sigAlg).build(kp.getPrivate())));
}
// --- PEM file helpers ---
private String toPem(Object obj) throws Exception {
java.io.StringWriter sw = new java.io.StringWriter();
try (JcaPEMWriter w = new JcaPEMWriter(sw)) {
w.writeObject(obj);
}
return sw.toString();
}
private Path writeTempPem(String prefix, Object... objects) throws Exception {
Path p = Files.createTempFile(prefix + "-", ".pem");
tempFiles.add(p);
try (JcaPEMWriter w = new JcaPEMWriter(Files.newBufferedWriter(p))) {
for (Object o : objects) {
w.writeObject(toPkcs8IfKey(o));
}
}
return p;
}
private Path writeTempPemEncrypted(String prefix, String password, Object... objects) throws Exception {
Path p = Files.createTempFile(prefix + "-", ".pem");
tempFiles.add(p);
var encryptor = new JcePEMEncryptorBuilder("AES-256-CBC")
.setProvider(BouncyCastleProvider.PROVIDER_NAME)
.build(password.toCharArray());
try (JcaPEMWriter w = new JcaPEMWriter(Files.newBufferedWriter(p))) {
for (Object o : objects) {
if (o instanceof PrivateKey) {
w.writeObject(o, encryptor);
} else {
w.writeObject(o);
}
}
}
return p;
}
private Object toPkcs8IfKey(Object o) {
if (o instanceof PrivateKey pk) {
return new PemObject("PRIVATE KEY", pk.getEncoded());
}
return o;
}
}

8
application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java

@ -115,6 +115,8 @@ import java.util.Map;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@ -1749,13 +1751,13 @@ public class EntityServiceTest extends AbstractControllerTest {
}
@Test
public void testFindTenantTelemetry() {
public void testFindTenantTelemetry() throws ExecutionException, InterruptedException, TimeoutException {
// save timeseries by sys admin
BasicTsKvEntry timeseries = new BasicTsKvEntry(42L, new DoubleDataEntry("temperature", 45.5));
timeseriesService.save(TenantId.SYS_TENANT_ID, tenantId, timeseries);
timeseriesService.save(TenantId.SYS_TENANT_ID, tenantId, timeseries).get(TIMEOUT, TimeUnit.SECONDS);
AttributeKvEntry attr = new BaseAttributeKvEntry(new LongDataEntry("attr", 10L), 42L);
attributesService.save(TenantId.SYS_TENANT_ID, tenantId, SERVER_SCOPE, List.of(attr));
attributesService.save(TenantId.SYS_TENANT_ID, tenantId, SERVER_SCOPE, List.of(attr)).get(TIMEOUT, TimeUnit.SECONDS);
SingleEntityFilter singleEntityFilter = new SingleEntityFilter();
singleEntityFilter.setSingleEntity(AliasEntityId.fromEntityId(tenantId));

4
common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.notification.rule;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@ -29,7 +30,7 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
@JsonIgnoreProperties
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "triggerType", visible = true, include = JsonTypeInfo.As.EXISTING_PROPERTY, defaultImpl = DefaultNotificationRuleRecipientsConfig.class)
@JsonSubTypes({
@Type(name = "ALARM", value = EscalatedNotificationRuleRecipientsConfig.class),
@ -38,6 +39,7 @@ import java.util.UUID;
public abstract class NotificationRuleRecipientsConfig implements Serializable {
@NotNull
@JsonProperty("triggerType")
private NotificationRuleTriggerType triggerType;
@JsonIgnore

1
common/edge-api/src/main/proto/edge.proto

@ -47,6 +47,7 @@ enum EdgeVersion {
V_4_3_0 = 13;
V_4_2_1_2 = 14;
V_4_2_2 = 4220;
V_4_2_2_1 = 4221;
V_4_3_0_1 = 15;
V_4_3_1 = 4310;

2
dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java

@ -65,7 +65,7 @@ public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest {
new BasicTsKvEntry(TimeUnit.MINUTES.toMillis(5), new JsonDataEntry("test", "{\"test\":\"testValue\"}")));
DeviceId deviceId = new DeviceId(Uuids.timeBased());
tsService.save(tenantId, deviceId, timeseries, ttlInSec);
tsService.save(tenantId, deviceId, timeseries, ttlInSec).get(MAX_TIMEOUT, TimeUnit.SECONDS);
List<TsKvEntry> fullList = tsService.findAll(tenantId, deviceId, Collections.singletonList(new BaseReadTsKvQuery("test", 0L,
TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE))).get(MAX_TIMEOUT, TimeUnit.SECONDS);

Loading…
Cancel
Save