diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index 76070bec53..4e7c0dad0c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/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. + *

+ * Delegates PEM parsing and key management to {@link PemSslCredentials} — the same + * class used by MQTT, CoAP, and LwM2M transports — which supports: + *

+ * 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) { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index f0f2e1eaa7..fc8c5c3be2 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/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()); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 2d1fc1b63f..73cdcb484e 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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. diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java index a4d6b62a0c..277ec687f1 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityViewControllerTest.java +++ b/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); diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index 87ba0ec3e8..eab878854d 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/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 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 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 tsData) throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); tsService.saveTimeseries(TimeseriesSaveRequest.builder() diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSslTest.java new file mode 100644 index 0000000000..04a1b87833 --- /dev/null +++ b/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. + *

+ * 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 + *

+ * 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 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; + } +} diff --git a/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java b/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java index c54ffe4a07..05ab1165f6 100644 --- a/application/src/test/java/org/thingsboard/server/service/entitiy/EntityServiceTest.java +++ b/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)); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java index fda73b8059..de4e05a2cb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/notification/rule/NotificationRuleRecipientsConfig.java +++ b/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 diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 913a9199e7..75f7c4190a 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/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; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java index aa43826734..b66cdb4363 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java +++ b/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 fullList = tsService.findAll(tenantId, deviceId, Collections.singletonList(new BaseReadTsKvQuery("test", 0L, TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE))).get(MAX_TIMEOUT, TimeUnit.SECONDS);