Browse Source
# Conflicts: # common/coap-server/src/test/java/org/thingsboard/server/coapserver/TbCoapDtlsSettingsTest.javapull/15487/head
275 changed files with 10215 additions and 1555 deletions
File diff suppressed because one or more lines are too long
@ -0,0 +1,141 @@ |
|||||
|
/** |
||||
|
* 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.entitiy.queue; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.cluster.TbClusterService; |
||||
|
import org.thingsboard.server.common.data.id.QueueId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.queue.Queue; |
||||
|
import org.thingsboard.server.dao.queue.QueueService; |
||||
|
import org.thingsboard.server.queue.TbQueueAdmin; |
||||
|
import org.thingsboard.server.queue.discovery.TopicService; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.anyBoolean; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultTbQueueServiceTest { |
||||
|
|
||||
|
@Mock |
||||
|
private QueueService queueServiceMock; |
||||
|
@Mock |
||||
|
private TbClusterService tbClusterServiceMock; |
||||
|
@Mock |
||||
|
private TbQueueAdmin tbQueueAdminMock; |
||||
|
|
||||
|
private TopicService topicService; |
||||
|
private DefaultTbQueueService tbQueueService; |
||||
|
|
||||
|
private final TenantId tenantId = TenantId.SYS_TENANT_ID; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
topicService = new TopicService(); |
||||
|
tbQueueService = new DefaultTbQueueService(queueServiceMock, tbClusterServiceMock, tbQueueAdminMock, topicService); |
||||
|
} |
||||
|
|
||||
|
private Queue newQueue(int partitions) { |
||||
|
Queue queue = new Queue(); |
||||
|
queue.setTenantId(tenantId); |
||||
|
queue.setName("testQueue"); |
||||
|
queue.setTopic("tb_rule_engine.testQueue"); |
||||
|
queue.setPartitions(partitions); |
||||
|
return queue; |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenQueuePrefix_whenSaveQueue_thenCreatesPrefixedTopics() { |
||||
|
// queue.prefix = "thingsboard" (TB_QUEUE_PREFIX set)
|
||||
|
ReflectionTestUtils.setField(topicService, "prefix", "thingsboard"); |
||||
|
|
||||
|
Queue queue = newQueue(2); |
||||
|
when(queueServiceMock.saveQueue(queue)).thenReturn(queue); |
||||
|
|
||||
|
tbQueueService.saveQueue(queue); |
||||
|
|
||||
|
ArgumentCaptor<String> topicCaptor = ArgumentCaptor.forClass(String.class); |
||||
|
verify(tbQueueAdminMock, times(2)).createTopicIfNotExists(topicCaptor.capture(), any(), anyBoolean()); |
||||
|
|
||||
|
// All created topics must carry the prefix - this is the fix.
|
||||
|
assertThat(topicCaptor.getAllValues()) |
||||
|
.containsExactlyInAnyOrder( |
||||
|
"thingsboard.tb_rule_engine.testQueue.0", |
||||
|
"thingsboard.tb_rule_engine.testQueue.1"); |
||||
|
// No unprefixed (orphan-prone) topic must ever be created.
|
||||
|
assertThat(topicCaptor.getAllValues()) |
||||
|
.noneMatch(topic -> topic.equals("tb_rule_engine.testQueue.0") |
||||
|
|| topic.equals("tb_rule_engine.testQueue.1")); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenNoQueuePrefix_whenSaveQueue_thenCreatesUnprefixedTopics() { |
||||
|
// queue.prefix blank (TB_QUEUE_PREFIX not set) - default behavior preserved
|
||||
|
ReflectionTestUtils.setField(topicService, "prefix", ""); |
||||
|
|
||||
|
Queue queue = newQueue(2); |
||||
|
when(queueServiceMock.saveQueue(queue)).thenReturn(queue); |
||||
|
|
||||
|
tbQueueService.saveQueue(queue); |
||||
|
|
||||
|
ArgumentCaptor<String> topicCaptor = ArgumentCaptor.forClass(String.class); |
||||
|
verify(tbQueueAdminMock, times(2)).createTopicIfNotExists(topicCaptor.capture(), any(), anyBoolean()); |
||||
|
|
||||
|
assertThat(topicCaptor.getAllValues()) |
||||
|
.containsExactlyInAnyOrder( |
||||
|
"tb_rule_engine.testQueue.0", |
||||
|
"tb_rule_engine.testQueue.1"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenQueuePrefix_whenIncreasePartitions_thenOnlyNewPartitionsCreatedPrefixed() { |
||||
|
ReflectionTestUtils.setField(topicService, "prefix", "thingsboard"); |
||||
|
|
||||
|
Queue oldQueue = newQueue(2); |
||||
|
oldQueue.setId(new QueueId(UUID.randomUUID())); |
||||
|
Queue updatedQueue = newQueue(4); |
||||
|
updatedQueue.setId(oldQueue.getId()); |
||||
|
|
||||
|
when(queueServiceMock.findQueueById(tenantId, updatedQueue.getId())).thenReturn(oldQueue); |
||||
|
when(queueServiceMock.saveQueue(updatedQueue)).thenReturn(updatedQueue); |
||||
|
|
||||
|
tbQueueService.saveQueue(updatedQueue); |
||||
|
|
||||
|
ArgumentCaptor<String> topicCaptor = ArgumentCaptor.forClass(String.class); |
||||
|
verify(tbQueueAdminMock, times(2)).createTopicIfNotExists(topicCaptor.capture(), any(), anyBoolean()); |
||||
|
|
||||
|
assertThat(topicCaptor.getAllValues()) |
||||
|
.containsExactlyInAnyOrder( |
||||
|
"thingsboard.tb_rule_engine.testQueue.2", |
||||
|
"thingsboard.tb_rule_engine.testQueue.3"); |
||||
|
verify(tbQueueAdminMock, never()).createTopicIfNotExists(eq("thingsboard.tb_rule_engine.testQueue.0"), any(), anyBoolean()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,78 @@ |
|||||
|
/** |
||||
|
* 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.queue; |
||||
|
|
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.data.rpc.RpcError; |
||||
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
||||
|
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; |
||||
|
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
||||
|
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.BDDMockito.then; |
||||
|
import static org.mockito.Mockito.doCallRealMethod; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultTbRuleEngineConsumerServiceTest { |
||||
|
|
||||
|
@Mock |
||||
|
private TbRuleEngineDeviceRpcService tbDeviceRpcServiceMock; |
||||
|
@Mock |
||||
|
private TbCallback tbCallbackMock; |
||||
|
|
||||
|
@Mock |
||||
|
private DefaultTbRuleEngineConsumerService defaultTbRuleEngineConsumerServiceMock; |
||||
|
|
||||
|
@Test |
||||
|
public void givenNotFoundErrorAndNoResponse_whenHandleFromDeviceRpcResponse_thenNotFoundAndNullResponseAreRecovered() { |
||||
|
// GIVEN
|
||||
|
ReflectionTestUtils.setField(defaultTbRuleEngineConsumerServiceMock, "tbDeviceRpcService", tbDeviceRpcServiceMock); |
||||
|
var requestId = UUID.randomUUID(); |
||||
|
// error = NOT_FOUND.ordinal() (0) and response left unset: the previously broken combination
|
||||
|
// ('error > 0' dropped NOT_FOUND, proto3 default collapsed a null response to "").
|
||||
|
var proto = TransportProtos.FromDeviceRPCResponseProto.newBuilder() |
||||
|
.setRequestIdMSB(requestId.getMostSignificantBits()) |
||||
|
.setRequestIdLSB(requestId.getLeastSignificantBits()) |
||||
|
.setError(RpcError.NOT_FOUND.ordinal()) |
||||
|
.build(); |
||||
|
var nfMsg = ToRuleEngineNotificationMsg.newBuilder().setFromDeviceRpcResponse(proto).build(); |
||||
|
var queueMsg = new TbProtoQueueMsg<>(requestId, nfMsg); |
||||
|
doCallRealMethod().when(defaultTbRuleEngineConsumerServiceMock).handleNotification(requestId, queueMsg, tbCallbackMock); |
||||
|
|
||||
|
// WHEN
|
||||
|
defaultTbRuleEngineConsumerServiceMock.handleNotification(requestId, queueMsg, tbCallbackMock); |
||||
|
|
||||
|
// THEN
|
||||
|
var responseCaptor = ArgumentCaptor.forClass(FromDeviceRpcResponse.class); |
||||
|
then(tbDeviceRpcServiceMock).should().processRpcResponseFromDevice(responseCaptor.capture()); |
||||
|
var response = responseCaptor.getValue(); |
||||
|
assertThat(response.getId()).isEqualTo(requestId); |
||||
|
assertThat(response.getError()).contains(RpcError.NOT_FOUND); |
||||
|
assertThat(response.getResponse()).isEmpty(); |
||||
|
then(tbCallbackMock).should().onSuccess(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,349 @@ |
|||||
|
/** |
||||
|
* 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.coapserver; |
||||
|
|
||||
|
import org.bouncycastle.asn1.x500.X500Name; |
||||
|
import org.bouncycastle.cert.jcajce.JcaX509CertificateConverter; |
||||
|
import org.bouncycastle.cert.jcajce.JcaX509v3CertificateBuilder; |
||||
|
import org.bouncycastle.operator.jcajce.JcaContentSignerBuilder; |
||||
|
import org.bouncycastle.util.io.pem.PemObject; |
||||
|
import org.bouncycastle.util.io.pem.PemWriter; |
||||
|
import org.eclipse.californium.core.CoapClient; |
||||
|
import org.eclipse.californium.core.CoapResource; |
||||
|
import org.eclipse.californium.core.CoapResponse; |
||||
|
import org.eclipse.californium.core.CoapServer; |
||||
|
import org.eclipse.californium.core.coap.CoAP; |
||||
|
import org.eclipse.californium.core.config.CoapConfig; |
||||
|
import org.eclipse.californium.core.network.CoapEndpoint; |
||||
|
import org.eclipse.californium.core.server.resources.CoapExchange; |
||||
|
import org.eclipse.californium.elements.config.Configuration; |
||||
|
import org.eclipse.californium.elements.util.SslContextUtil; |
||||
|
import org.eclipse.californium.scandium.DTLSConnector; |
||||
|
import org.eclipse.californium.scandium.config.DtlsConfig; |
||||
|
import org.eclipse.californium.scandium.config.DtlsConnectorConfig; |
||||
|
import org.eclipse.californium.scandium.dtls.CertificateType; |
||||
|
import org.eclipse.californium.scandium.dtls.x509.SingleCertificateProvider; |
||||
|
import org.eclipse.californium.scandium.dtls.x509.StaticNewAdvancedCertificateVerifier; |
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.io.TempDir; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.KeystoreSslCredentials; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.PemSslCredentials; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.SslCredentials; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.SslCredentialsConfig; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.SslCredentialsType; |
||||
|
|
||||
|
import java.io.OutputStreamWriter; |
||||
|
import java.math.BigInteger; |
||||
|
import java.net.InetAddress; |
||||
|
import java.net.InetSocketAddress; |
||||
|
import java.nio.file.Files; |
||||
|
import java.nio.file.Path; |
||||
|
import java.security.KeyPair; |
||||
|
import java.security.KeyPairGenerator; |
||||
|
import java.security.cert.X509Certificate; |
||||
|
import java.util.Collections; |
||||
|
import java.util.Date; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static java.util.concurrent.TimeUnit.MILLISECONDS; |
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_CLIENT_AUTHENTICATION_MODE; |
||||
|
import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RETRANSMISSION_TIMEOUT; |
||||
|
import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_ROLE; |
||||
|
import static org.eclipse.californium.scandium.config.DtlsConfig.DtlsRole.SERVER_ONLY; |
||||
|
|
||||
|
public class CoapDtlsCertificateReloadIntegrationTest { |
||||
|
|
||||
|
private static final String TEST_RESOURCE_PATH = "test"; |
||||
|
private static final String TEST_PAYLOAD = "hello-dtls"; |
||||
|
|
||||
|
@TempDir |
||||
|
Path tempDir; |
||||
|
|
||||
|
private CoapServer coapServer; |
||||
|
|
||||
|
@AfterEach |
||||
|
public void teardown() { |
||||
|
if (coapServer != null) { |
||||
|
coapServer.destroy(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsServer_whenCertFileChangedAndReloadTriggered_thenNewEndpointServesNewCert() throws Exception { |
||||
|
KeyPair keyPairA = generateKeyPair(); |
||||
|
X509Certificate certA = generateSelfSignedCert(keyPairA, "CN=ServerA"); |
||||
|
KeyPair keyPairB = generateKeyPair(); |
||||
|
X509Certificate certB = generateSelfSignedCert(keyPairB, "CN=ServerB"); |
||||
|
|
||||
|
Path certFile = tempDir.resolve("server-cert.pem"); |
||||
|
Path keyFile = tempDir.resolve("server-key.pem"); |
||||
|
writeCertPem(certFile, certA); |
||||
|
writeKeyPem(keyFile, keyPairA); |
||||
|
|
||||
|
SslCredentialsConfig credentialsConfig = createSslCredentialsConfig(certFile, keyFile); |
||||
|
|
||||
|
Configuration config = createServerConfig(); |
||||
|
coapServer = new CoapServer(config); |
||||
|
coapServer.add(new TestResource()); |
||||
|
|
||||
|
int dtlsPort = findAvailablePort(); |
||||
|
CoapEndpoint endpointA = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpointA); |
||||
|
coapServer.start(); |
||||
|
|
||||
|
CoapResponse responseA = doDtlsRequest(dtlsPort, certA); |
||||
|
assertThat(responseA).isNotNull(); |
||||
|
assertThat(responseA.getCode()).isEqualTo(CoAP.ResponseCode.CONTENT); |
||||
|
assertThat(responseA.getResponseText()).isEqualTo(TEST_PAYLOAD); |
||||
|
|
||||
|
writeCertPem(certFile, certB); |
||||
|
writeKeyPem(keyFile, keyPairB); |
||||
|
credentialsConfig.onCertificateFileChanged(); |
||||
|
|
||||
|
coapServer.getEndpoints().remove(endpointA); |
||||
|
endpointA.stop(); |
||||
|
|
||||
|
CoapEndpoint endpointB = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpointB); |
||||
|
endpointB.start(); |
||||
|
endpointA.destroy(); |
||||
|
|
||||
|
CoapResponse responseB = doDtlsRequest(dtlsPort, certB); |
||||
|
assertThat(responseB).isNotNull(); |
||||
|
assertThat(responseB.getCode()).isEqualTo(CoAP.ResponseCode.CONTENT); |
||||
|
assertThat(responseB.getResponseText()).isEqualTo(TEST_PAYLOAD); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsServer_whenCertReloaded_thenOldCertClientFails() throws Exception { |
||||
|
KeyPair keyPairA = generateKeyPair(); |
||||
|
X509Certificate certA = generateSelfSignedCert(keyPairA, "CN=ServerA"); |
||||
|
KeyPair keyPairB = generateKeyPair(); |
||||
|
X509Certificate certB = generateSelfSignedCert(keyPairB, "CN=ServerB"); |
||||
|
|
||||
|
Path certFile = tempDir.resolve("server-cert.pem"); |
||||
|
Path keyFile = tempDir.resolve("server-key.pem"); |
||||
|
writeCertPem(certFile, certA); |
||||
|
writeKeyPem(keyFile, keyPairA); |
||||
|
|
||||
|
SslCredentialsConfig credentialsConfig = createSslCredentialsConfig(certFile, keyFile); |
||||
|
|
||||
|
Configuration config = createServerConfig(); |
||||
|
coapServer = new CoapServer(config); |
||||
|
coapServer.add(new TestResource()); |
||||
|
|
||||
|
int dtlsPort = findAvailablePort(); |
||||
|
CoapEndpoint endpointA = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpointA); |
||||
|
coapServer.start(); |
||||
|
|
||||
|
CoapResponse responseA = doDtlsRequest(dtlsPort, certA); |
||||
|
assertThat(responseA).isNotNull(); |
||||
|
|
||||
|
writeCertPem(certFile, certB); |
||||
|
writeKeyPem(keyFile, keyPairB); |
||||
|
credentialsConfig.onCertificateFileChanged(); |
||||
|
|
||||
|
coapServer.getEndpoints().remove(endpointA); |
||||
|
endpointA.stop(); |
||||
|
CoapEndpoint endpointB = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpointB); |
||||
|
endpointB.start(); |
||||
|
endpointA.destroy(); |
||||
|
|
||||
|
CoapResponse failedResponse = doDtlsRequest(dtlsPort, certA); |
||||
|
assertThat(failedResponse).isNull(); |
||||
|
|
||||
|
CoapResponse responseB = doDtlsRequest(dtlsPort, certB); |
||||
|
assertThat(responseB).isNotNull(); |
||||
|
assertThat(responseB.getCode()).isEqualTo(CoAP.ResponseCode.CONTENT); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsServer_whenReloadWithSameCert_thenConnectionStillWorks() throws Exception { |
||||
|
KeyPair keyPair = generateKeyPair(); |
||||
|
X509Certificate cert = generateSelfSignedCert(keyPair, "CN=Server"); |
||||
|
|
||||
|
Path certFile = tempDir.resolve("server-cert.pem"); |
||||
|
Path keyFile = tempDir.resolve("server-key.pem"); |
||||
|
writeCertPem(certFile, cert); |
||||
|
writeKeyPem(keyFile, keyPair); |
||||
|
|
||||
|
SslCredentialsConfig credentialsConfig = createSslCredentialsConfig(certFile, keyFile); |
||||
|
|
||||
|
Configuration config = createServerConfig(); |
||||
|
coapServer = new CoapServer(config); |
||||
|
coapServer.add(new TestResource()); |
||||
|
|
||||
|
int dtlsPort = findAvailablePort(); |
||||
|
CoapEndpoint endpoint1 = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpoint1); |
||||
|
coapServer.start(); |
||||
|
|
||||
|
CoapResponse response1 = doDtlsRequest(dtlsPort, cert); |
||||
|
assertThat(response1).isNotNull(); |
||||
|
assertThat(response1.getCode()).isEqualTo(CoAP.ResponseCode.CONTENT); |
||||
|
|
||||
|
credentialsConfig.onCertificateFileChanged(); |
||||
|
|
||||
|
coapServer.getEndpoints().remove(endpoint1); |
||||
|
endpoint1.stop(); |
||||
|
CoapEndpoint endpoint2 = buildDtlsEndpointFromCredentials(config, credentialsConfig.getCredentials(), dtlsPort); |
||||
|
coapServer.addEndpoint(endpoint2); |
||||
|
endpoint2.start(); |
||||
|
endpoint1.destroy(); |
||||
|
|
||||
|
CoapResponse response2 = doDtlsRequest(dtlsPort, cert); |
||||
|
assertThat(response2).isNotNull(); |
||||
|
assertThat(response2.getCode()).isEqualTo(CoAP.ResponseCode.CONTENT); |
||||
|
} |
||||
|
|
||||
|
private SslCredentialsConfig createSslCredentialsConfig(Path certFile, Path keyFile) { |
||||
|
PemSslCredentials pem = new PemSslCredentials(); |
||||
|
pem.setCertFile(certFile.toAbsolutePath().toString()); |
||||
|
pem.setKeyFile(keyFile.toAbsolutePath().toString()); |
||||
|
|
||||
|
SslCredentialsConfig config = new SslCredentialsConfig("CoAP DTLS Test", false); |
||||
|
config.setEnabled(true); |
||||
|
config.setType(SslCredentialsType.PEM); |
||||
|
config.setPem(pem); |
||||
|
config.setKeystore(new KeystoreSslCredentials()); |
||||
|
config.init(); |
||||
|
return config; |
||||
|
} |
||||
|
|
||||
|
private CoapEndpoint buildDtlsEndpointFromCredentials(Configuration config, SslCredentials credentials, int port) { |
||||
|
DtlsConnectorConfig.Builder dtlsBuilder = new DtlsConnectorConfig.Builder(config); |
||||
|
dtlsBuilder.setAddress(new InetSocketAddress(InetAddress.getLoopbackAddress(), port)); |
||||
|
dtlsBuilder.set(DTLS_ROLE, SERVER_ONLY); |
||||
|
dtlsBuilder.set(DTLS_RETRANSMISSION_TIMEOUT, 3000, MILLISECONDS); |
||||
|
dtlsBuilder.set(DTLS_CLIENT_AUTHENTICATION_MODE, |
||||
|
org.eclipse.californium.elements.config.CertificateAuthenticationMode.WANTED); |
||||
|
|
||||
|
SslContextUtil.Credentials serverCreds = new SslContextUtil.Credentials( |
||||
|
credentials.getPrivateKey(), null, credentials.getCertificateChain()); |
||||
|
|
||||
|
dtlsBuilder.setCertificateIdentityProvider( |
||||
|
new SingleCertificateProvider(serverCreds.getPrivateKey(), serverCreds.getCertificateChain(), |
||||
|
Collections.singletonList(CertificateType.X_509))); |
||||
|
|
||||
|
dtlsBuilder.setAdvancedCertificateVerifier( |
||||
|
StaticNewAdvancedCertificateVerifier.builder() |
||||
|
.setTrustAllCertificates() |
||||
|
.build()); |
||||
|
|
||||
|
DTLSConnector connector = new DTLSConnector(dtlsBuilder.build()); |
||||
|
|
||||
|
CoapEndpoint.Builder endpointBuilder = new CoapEndpoint.Builder(); |
||||
|
endpointBuilder.setConfiguration(config); |
||||
|
endpointBuilder.setConnector(connector); |
||||
|
return endpointBuilder.build(); |
||||
|
} |
||||
|
|
||||
|
private KeyPair generateKeyPair() throws Exception { |
||||
|
KeyPairGenerator kpg = KeyPairGenerator.getInstance("EC"); |
||||
|
kpg.initialize(256); |
||||
|
return kpg.generateKeyPair(); |
||||
|
} |
||||
|
|
||||
|
private X509Certificate generateSelfSignedCert(KeyPair kp, String subjectDn) throws Exception { |
||||
|
X500Name subject = new X500Name(subjectDn); |
||||
|
Date now = new Date(); |
||||
|
Date expiry = new Date(now.getTime() + TimeUnit.DAYS.toMillis(1)); |
||||
|
return new JcaX509CertificateConverter().getCertificate( |
||||
|
new JcaX509v3CertificateBuilder( |
||||
|
subject, BigInteger.valueOf(System.nanoTime()), now, expiry, |
||||
|
subject, kp.getPublic()) |
||||
|
.build(new JcaContentSignerBuilder("SHA256withECDSA").build(kp.getPrivate()))); |
||||
|
} |
||||
|
|
||||
|
private void writeCertPem(Path path, X509Certificate cert) throws Exception { |
||||
|
try (PemWriter writer = new PemWriter(new OutputStreamWriter(Files.newOutputStream(path)))) { |
||||
|
writer.writeObject(new PemObject("CERTIFICATE", cert.getEncoded())); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void writeKeyPem(Path path, KeyPair keyPair) throws Exception { |
||||
|
try (PemWriter writer = new PemWriter(new OutputStreamWriter(Files.newOutputStream(path)))) { |
||||
|
writer.writeObject(new PemObject("PRIVATE KEY", keyPair.getPrivate().getEncoded())); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private Configuration createServerConfig() { |
||||
|
Configuration config = new Configuration(); |
||||
|
config.set(CoapConfig.MAX_RETRANSMIT, 2); |
||||
|
config.set(CoapConfig.RESPONSE_MATCHING, CoapConfig.MatcherMode.RELAXED); |
||||
|
return config; |
||||
|
} |
||||
|
|
||||
|
private CoapResponse doDtlsRequest(int port, X509Certificate trustedCert) { |
||||
|
try { |
||||
|
Configuration clientConfig = new Configuration(); |
||||
|
clientConfig.set(CoapConfig.MAX_RETRANSMIT, 1); |
||||
|
clientConfig.set(DtlsConfig.DTLS_ROLE, DtlsConfig.DtlsRole.CLIENT_ONLY); |
||||
|
clientConfig.set(DtlsConfig.DTLS_RETRANSMISSION_TIMEOUT, 2000, MILLISECONDS); |
||||
|
clientConfig.set(DtlsConfig.DTLS_USE_HELLO_VERIFY_REQUEST, false); |
||||
|
clientConfig.set(DtlsConfig.DTLS_VERIFY_SERVER_CERTIFICATES_SUBJECT, false); |
||||
|
|
||||
|
DtlsConnectorConfig.Builder clientDtls = new DtlsConnectorConfig.Builder(clientConfig); |
||||
|
clientDtls.setAdvancedCertificateVerifier( |
||||
|
StaticNewAdvancedCertificateVerifier.builder() |
||||
|
.setTrustedCertificates(trustedCert) |
||||
|
.build()); |
||||
|
|
||||
|
DTLSConnector clientConnector = new DTLSConnector(clientDtls.build()); |
||||
|
CoapEndpoint clientEndpoint = new CoapEndpoint.Builder() |
||||
|
.setConfiguration(clientConfig) |
||||
|
.setConnector(clientConnector) |
||||
|
.build(); |
||||
|
|
||||
|
CoapClient client = new CoapClient("coaps://127.0.0.1:" + port + "/" + TEST_RESOURCE_PATH); |
||||
|
client.setEndpoint(clientEndpoint); |
||||
|
client.setTimeout((long) 5000); |
||||
|
|
||||
|
try { |
||||
|
clientEndpoint.start(); |
||||
|
return client.get(); |
||||
|
} finally { |
||||
|
client.shutdown(); |
||||
|
clientEndpoint.destroy(); |
||||
|
} |
||||
|
} catch (Exception e) { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private int findAvailablePort() throws Exception { |
||||
|
try (java.net.DatagramSocket socket = new java.net.DatagramSocket(0)) { |
||||
|
return socket.getLocalPort(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private static class TestResource extends CoapResource { |
||||
|
TestResource() { |
||||
|
super(TEST_RESOURCE_PATH); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void handleGET(CoapExchange exchange) { |
||||
|
exchange.respond(CoAP.ResponseCode.CONTENT, TEST_PAYLOAD); |
||||
|
} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,246 @@ |
|||||
|
/** |
||||
|
* 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.coapserver; |
||||
|
|
||||
|
import org.eclipse.californium.core.CoapServer; |
||||
|
import org.eclipse.californium.core.network.CoapEndpoint; |
||||
|
import org.eclipse.californium.core.network.Endpoint; |
||||
|
import org.eclipse.californium.scandium.DTLSConnector; |
||||
|
import org.eclipse.californium.scandium.config.DtlsConnectorConfig; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.MockedConstruction; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.mockito.junit.jupiter.MockitoSettings; |
||||
|
import org.mockito.quality.Strictness; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
|
||||
|
import java.io.IOException; |
||||
|
import java.net.InetSocketAddress; |
||||
|
import java.util.List; |
||||
|
import java.util.concurrent.CopyOnWriteArrayList; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.doAnswer; |
||||
|
import static org.mockito.Mockito.doThrow; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.mockConstruction; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
@MockitoSettings(strictness = Strictness.LENIENT) |
||||
|
public class CoapDtlsCertificateReloadTest { |
||||
|
|
||||
|
@Mock |
||||
|
private CoapServerContext mockCoapServerContext; |
||||
|
|
||||
|
@Mock |
||||
|
private TbCoapDtlsSettings mockDtlsSettings; |
||||
|
|
||||
|
@Mock |
||||
|
private CoapServer mockCoapServer; |
||||
|
|
||||
|
@Mock |
||||
|
private CoapEndpoint mockDtlsEndpoint; |
||||
|
|
||||
|
@Mock |
||||
|
private DTLSConnector mockDtlsConnector; |
||||
|
|
||||
|
private DefaultCoapServerService coapServerService; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setup() { |
||||
|
coapServerService = new DefaultCoapServerService(); |
||||
|
ReflectionTestUtils.setField(coapServerService, "coapServerContext", mockCoapServerContext); |
||||
|
|
||||
|
when(mockCoapServerContext.getHost()).thenReturn("localhost"); |
||||
|
when(mockCoapServerContext.getPort()).thenReturn(5683); |
||||
|
doAnswer(invocation -> { |
||||
|
invocation.getArgument(0); |
||||
|
return null; |
||||
|
}).when(mockDtlsSettings).registerReloadCallback(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsEnabled_whenRegisterCertificateReloadCallback_thenShouldRegisterCallback() { |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
|
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
|
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "afterSingletonsInstantiated"); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
verify(mockDtlsSettings).registerReloadCallback(callbackCaptor.capture()); |
||||
|
assertThat(callbackCaptor.getValue()).isNotNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsNotEnabled_whenRegisterCertificateReloadCallback_thenShouldNotRegisterCallback() { |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(null); |
||||
|
|
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "afterSingletonsInstantiated"); |
||||
|
|
||||
|
verify(mockDtlsSettings, never()).registerReloadCallback(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallbackInvoked_whenNewEndpointCreationFails_thenOldEndpointIsPreserved() { |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
|
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsCoapEndpoint", mockDtlsEndpoint); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsConnector", mockDtlsConnector); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "afterSingletonsInstantiated"); |
||||
|
verify(mockDtlsSettings).registerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
// dtlsSettings.dtlsConnectorConfig() isn't mocked, so the callback will throw.
|
||||
|
// The old endpoint should not be stopped/destroyed when creation of the new one fails.
|
||||
|
reloadCallback.run(); |
||||
|
|
||||
|
verify(mockDtlsEndpoint, never()).stop(); |
||||
|
verify(mockDtlsConnector, never()).destroy(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenDtlsEnabled_whenInit_thenShouldRegisterCallback() { |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
when(mockCoapServerContext.getHost()).thenReturn("localhost"); |
||||
|
when(mockCoapServerContext.getPort()).thenReturn(5683); |
||||
|
|
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "afterSingletonsInstantiated"); |
||||
|
|
||||
|
verify(mockDtlsSettings).registerReloadCallback(any(Runnable.class)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenInvokedMultipleTimes_thenShouldRegisterOnce() { |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsCoapEndpoint", mockDtlsEndpoint); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsConnector", mockDtlsConnector); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "afterSingletonsInstantiated"); |
||||
|
verify(mockDtlsSettings).registerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
assertThat(reloadCallback).isNotNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenSuccessful_thenOldEndpointRemovedFromServer() throws Exception { |
||||
|
// GIVEN
|
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
|
||||
|
DtlsConnectorConfig mockDtlsConfig = mock(DtlsConnectorConfig.class); |
||||
|
TbCoapDtlsCertificateVerifier mockNewVerifier = mock(TbCoapDtlsCertificateVerifier.class); |
||||
|
when(mockDtlsConfig.getAdvancedCertificateVerifier()).thenReturn(mockNewVerifier); |
||||
|
when(mockDtlsConfig.getAddress()).thenReturn(new InetSocketAddress("localhost", 5684)); |
||||
|
when(mockDtlsSettings.dtlsConnectorConfig(any())).thenReturn(mockDtlsConfig); |
||||
|
|
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsCoapEndpoint", mockDtlsEndpoint); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsConnector", mockDtlsConnector); |
||||
|
|
||||
|
List<Endpoint> endpointsList = new CopyOnWriteArrayList<>(); |
||||
|
endpointsList.add(mockDtlsEndpoint); |
||||
|
when(mockCoapServer.getEndpoints()).thenReturn(endpointsList); |
||||
|
|
||||
|
CoapEndpoint mockNewEndpoint = mock(CoapEndpoint.class); |
||||
|
|
||||
|
try (MockedConstruction<DTLSConnector> dtlsMock = mockConstruction(DTLSConnector.class); |
||||
|
MockedConstruction<CoapEndpoint.Builder> builderMock = mockConstruction(CoapEndpoint.Builder.class, |
||||
|
(builder, context) -> { |
||||
|
when(builder.build()).thenReturn(mockNewEndpoint); |
||||
|
when(builder.setConfiguration(any())).thenReturn(builder); |
||||
|
when(builder.setConnector(any(DTLSConnector.class))).thenReturn(builder); |
||||
|
})) { |
||||
|
|
||||
|
// WHEN
|
||||
|
ReflectionTestUtils.invokeMethod(coapServerService, "recreateDtlsEndpoint"); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(endpointsList).doesNotContain(mockDtlsEndpoint); |
||||
|
verify(mockDtlsEndpoint).stop(); |
||||
|
verify(mockDtlsEndpoint).destroy(); |
||||
|
verify(mockDtlsConnector).destroy(); |
||||
|
verify(mockCoapServer).addEndpoint(mockNewEndpoint); |
||||
|
verify(mockNewEndpoint).start(); |
||||
|
assertThat(ReflectionTestUtils.getField(coapServerService, "dtlsCoapEndpoint")).isSameAs(mockNewEndpoint); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenStartFails_thenNewResourcesCleanedAndOldRestored() throws Exception { |
||||
|
// GIVEN
|
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
|
||||
|
DtlsConnectorConfig mockDtlsConfig = mock(DtlsConnectorConfig.class); |
||||
|
when(mockDtlsConfig.getAddress()).thenReturn(new InetSocketAddress("localhost", 5684)); |
||||
|
when(mockDtlsSettings.dtlsConnectorConfig(any())).thenReturn(mockDtlsConfig); |
||||
|
|
||||
|
ReflectionTestUtils.setField(coapServerService, "server", mockCoapServer); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsCoapEndpoint", mockDtlsEndpoint); |
||||
|
ReflectionTestUtils.setField(coapServerService, "dtlsConnector", mockDtlsConnector); |
||||
|
|
||||
|
List<Endpoint> endpointsList = new CopyOnWriteArrayList<>(); |
||||
|
endpointsList.add(mockDtlsEndpoint); |
||||
|
when(mockCoapServer.getEndpoints()).thenReturn(endpointsList); |
||||
|
|
||||
|
CoapEndpoint mockNewEndpoint = mock(CoapEndpoint.class); |
||||
|
doThrow(new IOException("start failed")).when(mockNewEndpoint).start(); |
||||
|
|
||||
|
try (MockedConstruction<DTLSConnector> dtlsMock = mockConstruction(DTLSConnector.class); |
||||
|
MockedConstruction<CoapEndpoint.Builder> builderMock = mockConstruction(CoapEndpoint.Builder.class, |
||||
|
(builder, context) -> { |
||||
|
when(builder.build()).thenReturn(mockNewEndpoint); |
||||
|
when(builder.setConfiguration(any())).thenReturn(builder); |
||||
|
when(builder.setConnector(any(DTLSConnector.class))).thenReturn(builder); |
||||
|
})) { |
||||
|
|
||||
|
// WHEN
|
||||
|
coapServerService.afterSingletonsInstantiated(); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
verify(mockDtlsSettings).registerReloadCallback(callbackCaptor.capture()); |
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
reloadCallback.run(); |
||||
|
|
||||
|
// THEN - new resources cleaned up
|
||||
|
DTLSConnector constructedConnector = dtlsMock.constructed().get(0); |
||||
|
verify(mockNewEndpoint).destroy(); |
||||
|
verify(constructedConnector).destroy(); |
||||
|
assertThat(endpointsList).doesNotContain(mockNewEndpoint); |
||||
|
// Old endpoint was stopped to release port, then restored after new one failed
|
||||
|
verify(mockDtlsEndpoint).stop(); |
||||
|
verify(mockDtlsEndpoint).start(); |
||||
|
// Old fields preserved
|
||||
|
assertThat(ReflectionTestUtils.getField(coapServerService, "dtlsCoapEndpoint")).isSameAs(mockDtlsEndpoint); |
||||
|
assertThat(ReflectionTestUtils.getField(coapServerService, "dtlsConnector")).isSameAs(mockDtlsConnector); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,145 @@ |
|||||
|
/** |
||||
|
* 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.coapserver; |
||||
|
|
||||
|
import org.eclipse.californium.core.CoapServer; |
||||
|
import org.eclipse.californium.core.network.CoapEndpoint; |
||||
|
import org.eclipse.californium.core.server.resources.Resource; |
||||
|
import org.eclipse.californium.scandium.DTLSConnector; |
||||
|
import org.eclipse.californium.scandium.config.DtlsConnectorConfig; |
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.MockedConstruction; |
||||
|
import org.mockito.MockedStatic; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
|
||||
|
import java.net.DatagramSocket; |
||||
|
import java.net.InetAddress; |
||||
|
import java.net.InetSocketAddress; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.anyString; |
||||
|
import static org.mockito.Mockito.doThrow; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.mockConstruction; |
||||
|
import static org.mockito.Mockito.mockStatic; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultCoapServerServiceTest { |
||||
|
|
||||
|
private static final String HOST = "127.0.0.1"; |
||||
|
|
||||
|
@Mock |
||||
|
private CoapServerContext mockCoapServerContext; |
||||
|
|
||||
|
private DefaultCoapServerService service; |
||||
|
private DatagramSocket occupiedSocket; |
||||
|
private int occupiedPort; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() throws Exception { |
||||
|
occupiedSocket = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0)); |
||||
|
occupiedPort = occupiedSocket.getLocalPort(); |
||||
|
|
||||
|
service = new DefaultCoapServerService(); |
||||
|
ReflectionTestUtils.setField(service, "coapServerContext", mockCoapServerContext); |
||||
|
|
||||
|
when(mockCoapServerContext.getHost()).thenReturn(HOST); |
||||
|
when(mockCoapServerContext.getPort()).thenReturn(occupiedPort); |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(null); |
||||
|
} |
||||
|
|
||||
|
@AfterEach |
||||
|
public void tearDown() { |
||||
|
if (occupiedSocket != null && !occupiedSocket.isClosed()) { |
||||
|
occupiedSocket.close(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void whenPlainBindFails_thenInitThrowsAndReleasesCoapServer() { |
||||
|
assertThatThrownBy(() -> service.init()) |
||||
|
.isInstanceOf(IllegalStateException.class) |
||||
|
.hasMessageContaining("None of the server endpoints could be started"); |
||||
|
|
||||
|
assertThat(ReflectionTestUtils.getField(service, "server")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsSessionsExecutor")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsConnector")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsCoapEndpoint")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "tbDtlsCertificateVerifier")).isNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void whenDtlsEnabledAndStartFails_thenInitShutsDownDtlsExecutorAndReleasesCoapServer() throws Exception { |
||||
|
// DTLS enabled: the DTLS endpoint is created and dtlsSessionsExecutor is scheduled before server.start().
|
||||
|
// This exercises the catch's dtlsSessionsExecutor.shutdownNow() branch, which the plain-bind test does not.
|
||||
|
TbCoapDtlsSettings mockDtlsSettings = mock(TbCoapDtlsSettings.class); |
||||
|
when(mockCoapServerContext.getDtlsSettings()).thenReturn(mockDtlsSettings); |
||||
|
|
||||
|
DtlsConnectorConfig mockDtlsConfig = mock(DtlsConnectorConfig.class); |
||||
|
when(mockDtlsConfig.getAddress()).thenReturn(new InetSocketAddress(InetAddress.getByName(HOST), occupiedPort + 1)); |
||||
|
TbCoapDtlsCertificateVerifier mockVerifier = mock(TbCoapDtlsCertificateVerifier.class); |
||||
|
when(mockVerifier.getDtlsSessionReportTimeout()).thenReturn(1800000L); |
||||
|
when(mockDtlsConfig.getAdvancedCertificateVerifier()).thenReturn(mockVerifier); |
||||
|
when(mockDtlsSettings.dtlsConnectorConfig(any())).thenReturn(mockDtlsConfig); |
||||
|
|
||||
|
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class); |
||||
|
Resource mockRoot = mock(Resource.class); |
||||
|
|
||||
|
try (MockedStatic<ThingsBoardExecutors> executorsStatic = mockStatic(ThingsBoardExecutors.class); |
||||
|
MockedConstruction<CoapServer> serverMock = mockConstruction(CoapServer.class, (server, ctx) -> { |
||||
|
when(server.getRoot()).thenReturn(mockRoot); |
||||
|
doThrow(new IllegalStateException("None of the server endpoints could be started")).when(server).start(); |
||||
|
}); |
||||
|
MockedConstruction<DTLSConnector> dtlsMock = mockConstruction(DTLSConnector.class); |
||||
|
MockedConstruction<CoapEndpoint.Builder> builderMock = mockConstruction(CoapEndpoint.Builder.class, (builder, ctx) -> { |
||||
|
when(builder.setInetSocketAddress(any())).thenReturn(builder); |
||||
|
when(builder.setConfiguration(any())).thenReturn(builder); |
||||
|
when(builder.setConnector(any(DTLSConnector.class))).thenReturn(builder); |
||||
|
when(builder.build()).thenReturn(mock(CoapEndpoint.class)); |
||||
|
})) { |
||||
|
|
||||
|
executorsStatic.when(() -> ThingsBoardExecutors.newSingleThreadScheduledExecutor(anyString())).thenReturn(mockExecutor); |
||||
|
|
||||
|
assertThatThrownBy(() -> service.init()) |
||||
|
.isInstanceOf(IllegalStateException.class) |
||||
|
.hasMessageContaining("None of the server endpoints could be started"); |
||||
|
|
||||
|
// DTLS branch was actually entered and the executor was created...
|
||||
|
verify(mockDtlsSettings).dtlsConnectorConfig(any()); |
||||
|
// ...and the cleanup branch shut it down and destroyed the server.
|
||||
|
verify(mockExecutor).shutdownNow(); |
||||
|
verify(serverMock.constructed().get(0)).destroy(); |
||||
|
} |
||||
|
|
||||
|
assertThat(ReflectionTestUtils.getField(service, "server")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsSessionsExecutor")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsConnector")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "dtlsCoapEndpoint")).isNull(); |
||||
|
assertThat(ReflectionTestUtils.getField(service, "tbDtlsCertificateVerifier")).isNull(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* 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.common.data; |
||||
|
|
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
|
||||
|
class ResourceUtilsTest { |
||||
|
|
||||
|
@Test |
||||
|
public void givenNonExistentResource_whenGetUri_thenThrowsRuntimeException() { |
||||
|
assertThatThrownBy(() -> ResourceUtils.getUri(ResourceUtilsTest.class.getClassLoader(), "non/existent/resource/path.txt")) |
||||
|
.isInstanceOf(RuntimeException.class) |
||||
|
.hasMessageContaining("Unable to find resource"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenExistingClasspathResource_whenGetUri_thenReturnsNonNullUri() { |
||||
|
String result = ResourceUtils.getUri(ResourceUtilsTest.class.getClassLoader(), "org/thingsboard/server/common/data/ResourceUtilsTest.class"); |
||||
|
|
||||
|
assertThat(result).isNotNull(); |
||||
|
assertThat(result).contains("ResourceUtilsTest"); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,146 @@ |
|||||
|
/** |
||||
|
* 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.queue.common; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
||||
|
import org.thingsboard.server.queue.TbQueueMsg; |
||||
|
|
||||
|
import java.util.Collections; |
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static org.hamcrest.MatcherAssert.assertThat; |
||||
|
import static org.hamcrest.Matchers.empty; |
||||
|
import static org.hamcrest.Matchers.greaterThanOrEqualTo; |
||||
|
import static org.hamcrest.Matchers.is; |
||||
|
import static org.mockito.ArgumentMatchers.anyLong; |
||||
|
import static org.mockito.BDDMockito.never; |
||||
|
import static org.mockito.BDDMockito.spy; |
||||
|
import static org.mockito.BDDMockito.times; |
||||
|
import static org.mockito.BDDMockito.verify; |
||||
|
|
||||
|
@Slf4j |
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class AbstractTbQueueConsumerTemplateTest { |
||||
|
|
||||
|
private static final long POLL_DURATION_MS = 100L; |
||||
|
private static final long SLEEP_TOLERANCE_MS = 20L; |
||||
|
|
||||
|
@Test |
||||
|
public void givenEmptyPartitionsAndLongPollingSupported_whenPoll_thenSleepsAndDoesNotCallDoPoll() { |
||||
|
// Regression: with empty partitions AND isLongPollingSupported()==true (e.g. Kafka),
|
||||
|
// poll() previously returned instantly with no sleep, causing the consumer loop to busy-spin.
|
||||
|
TestConsumer consumer = spy(new TestConsumer("test-topic", true)); |
||||
|
consumer.subscribe(Collections.emptySet()); |
||||
|
|
||||
|
long startNs = System.nanoTime(); |
||||
|
List<TbQueueMsg> result = consumer.poll(POLL_DURATION_MS); |
||||
|
long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs); |
||||
|
|
||||
|
assertThat(result, is(empty())); |
||||
|
verify(consumer, never()).doPoll(anyLong()); |
||||
|
assertThat("poll() must sleep ~durationInMillis when partitions are empty (no busy-wait)", |
||||
|
elapsedMs, greaterThanOrEqualTo(POLL_DURATION_MS - SLEEP_TOLERANCE_MS)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenEmptyPartitionsAndNoLongPolling_whenPoll_thenSleepsAndDoesNotCallDoPoll() { |
||||
|
TestConsumer consumer = spy(new TestConsumer("test-topic", false)); |
||||
|
consumer.subscribe(Collections.emptySet()); |
||||
|
|
||||
|
long startNs = System.nanoTime(); |
||||
|
List<TbQueueMsg> result = consumer.poll(POLL_DURATION_MS); |
||||
|
long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs); |
||||
|
|
||||
|
assertThat(result, is(empty())); |
||||
|
verify(consumer, never()).doPoll(anyLong()); |
||||
|
assertThat(elapsedMs, greaterThanOrEqualTo(POLL_DURATION_MS - SLEEP_TOLERANCE_MS)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenNonEmptyPartitions_whenPoll_thenCallsDoPoll() { |
||||
|
TestConsumer consumer = spy(new TestConsumer("test-topic", true)); |
||||
|
consumer.subscribe(Collections.singleton(new TopicPartitionInfo("test-topic", null, 0, true))); |
||||
|
|
||||
|
List<TbQueueMsg> result = consumer.poll(POLL_DURATION_MS); |
||||
|
|
||||
|
assertThat(result, is(empty())); |
||||
|
verify(consumer, times(1)).doPoll(POLL_DURATION_MS); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenPartitionsBecomeEmptyAfterRebalance_whenPollAgain_thenStopsCallingDoPoll() { |
||||
|
// Reproduces the observed trigger: a rebalance leaves the consumer with an empty
|
||||
|
// partition assignment. Subsequent poll() calls must not busy-spin or call doPoll().
|
||||
|
TestConsumer consumer = spy(new TestConsumer("test-topic", true)); |
||||
|
consumer.subscribe(Collections.singleton(new TopicPartitionInfo("test-topic", null, 0, true))); |
||||
|
consumer.poll(POLL_DURATION_MS); |
||||
|
verify(consumer, times(1)).doPoll(POLL_DURATION_MS); |
||||
|
|
||||
|
consumer.subscribe(Collections.emptySet()); |
||||
|
|
||||
|
long startNs = System.nanoTime(); |
||||
|
List<TbQueueMsg> result = consumer.poll(POLL_DURATION_MS); |
||||
|
long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs); |
||||
|
|
||||
|
assertThat(result, is(empty())); |
||||
|
verify(consumer, times(1)).doPoll(anyLong()); |
||||
|
assertThat(elapsedMs, greaterThanOrEqualTo(POLL_DURATION_MS - SLEEP_TOLERANCE_MS)); |
||||
|
} |
||||
|
|
||||
|
static class TestConsumer extends AbstractTbQueueConsumerTemplate<Object, TbQueueMsg> { |
||||
|
|
||||
|
private final boolean longPollingSupported; |
||||
|
|
||||
|
TestConsumer(String topic, boolean longPollingSupported) { |
||||
|
super(topic); |
||||
|
this.longPollingSupported = longPollingSupported; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected List<Object> doPoll(long durationInMillis) { |
||||
|
return Collections.emptyList(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected TbQueueMsg decode(Object record) { |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void doSubscribe(Set<TopicPartitionInfo> partitions) { |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void doCommit() { |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void doUnsubscribe() { |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected boolean isLongPollingSupported() { |
||||
|
return longPollingSupported; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,160 @@ |
|||||
|
/** |
||||
|
* 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.queue.notification; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.springframework.cache.Cache; |
||||
|
import org.springframework.cache.CacheManager; |
||||
|
import org.springframework.cache.concurrent.ConcurrentMapCacheManager; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.data.CacheConstants; |
||||
|
import org.thingsboard.server.common.data.notification.rule.NotificationRule; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTrigger; |
||||
|
import org.thingsboard.server.common.data.notification.rule.trigger.config.NotificationRuleTriggerType; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.concurrent.CopyOnWriteArrayList; |
||||
|
import java.util.concurrent.CyclicBarrier; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
class DefaultNotificationDeduplicationServiceTest { |
||||
|
|
||||
|
private static final int TIMEOUT = 30; |
||||
|
|
||||
|
private DefaultNotificationDeduplicationService deduplicationService; |
||||
|
private CacheManager cacheManager; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
deduplicationService = new DefaultNotificationDeduplicationService(); |
||||
|
deduplicationService.setDeduplicationDurations(""); |
||||
|
cacheManager = new ConcurrentMapCacheManager(CacheConstants.SENT_NOTIFICATIONS_CACHE); |
||||
|
ReflectionTestUtils.setField(deduplicationService, "cacheManager", cacheManager); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testFirstTriggerIsNotDeduplicated() { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(TimeUnit.HOURS.toMillis(1)); |
||||
|
NotificationRule rule = mockRule(); |
||||
|
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testSecondTriggerIsDeduplicated() { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(TimeUnit.HOURS.toMillis(1)); |
||||
|
NotificationRule rule = mockRule(); |
||||
|
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isFalse(); |
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testTriggerPassesAfterDeduplicationWindowExpires() { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(50); // 50ms dedup window
|
||||
|
NotificationRule rule = mockRule(); |
||||
|
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isFalse(); |
||||
|
|
||||
|
try { |
||||
|
Thread.sleep(200); // wait well past the 50ms window
|
||||
|
} catch (InterruptedException ignored) {} |
||||
|
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testFutureTimestampFromExternalCacheIsDiscarded() { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(TimeUnit.HOURS.toMillis(1)); |
||||
|
NotificationRule rule = mockRule(); |
||||
|
String dedupKey = DefaultNotificationDeduplicationService.getDeduplicationKey(trigger, rule); |
||||
|
|
||||
|
// Put a timestamp 2 hours in the future into external cache
|
||||
|
Cache externalCache = cacheManager.getCache(CacheConstants.SENT_NOTIFICATIONS_CACHE); |
||||
|
externalCache.put(dedupKey, System.currentTimeMillis() + TimeUnit.HOURS.toMillis(2)); |
||||
|
|
||||
|
// Should NOT be deduplicated — future timestamp must be discarded
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testValidTimestampFromExternalCacheIsDeduplicated() { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(TimeUnit.HOURS.toMillis(1)); |
||||
|
NotificationRule rule = mockRule(); |
||||
|
String dedupKey = DefaultNotificationDeduplicationService.getDeduplicationKey(trigger, rule); |
||||
|
|
||||
|
// Put a recent timestamp into external cache
|
||||
|
Cache externalCache = cacheManager.getCache(CacheConstants.SENT_NOTIFICATIONS_CACHE); |
||||
|
externalCache.put(dedupKey, System.currentTimeMillis()); |
||||
|
|
||||
|
// Should be deduplicated — valid external cache entry
|
||||
|
assertThat(deduplicationService.alreadyProcessed(trigger, rule)).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testConcurrentTriggersProduceExactlyOneNonDeduplicated() throws Exception { |
||||
|
NotificationRuleTrigger trigger = mockTrigger(TimeUnit.HOURS.toMillis(1)); |
||||
|
NotificationRule rule = mockRule(); |
||||
|
|
||||
|
int threadCount = 10; |
||||
|
CyclicBarrier barrier = new CyclicBarrier(threadCount); |
||||
|
List<Boolean> results = new CopyOnWriteArrayList<>(); |
||||
|
|
||||
|
ExecutorService executor = Executors.newFixedThreadPool(threadCount); |
||||
|
try { |
||||
|
for (int i = 0; i < threadCount; i++) { |
||||
|
executor.submit(() -> { |
||||
|
try { |
||||
|
barrier.await(TIMEOUT, TimeUnit.SECONDS); |
||||
|
} catch (Exception ignored) {} |
||||
|
results.add(deduplicationService.alreadyProcessed(trigger, rule)); |
||||
|
}); |
||||
|
} |
||||
|
executor.shutdown(); |
||||
|
assertThat(executor.awaitTermination(TIMEOUT, TimeUnit.SECONDS)).isTrue(); |
||||
|
|
||||
|
assertThat(results).hasSize(threadCount); |
||||
|
assertThat(results.stream().filter(r -> !r).count()) |
||||
|
.as("exactly one trigger should pass through deduplication") |
||||
|
.isEqualTo(1); |
||||
|
} finally { |
||||
|
executor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private NotificationRuleTrigger mockTrigger(long deduplicationDurationMs) { |
||||
|
NotificationRuleTrigger trigger = mock(NotificationRuleTrigger.class); |
||||
|
when(trigger.getType()).thenReturn(NotificationRuleTriggerType.RESOURCES_SHORTAGE); |
||||
|
when(trigger.getDeduplicationKey()).thenReturn("test:dedup:key"); |
||||
|
when(trigger.getDefaultDeduplicationDuration()).thenReturn(deduplicationDurationMs); |
||||
|
when(trigger.getDeduplicationStrategy()).thenReturn(NotificationRuleTrigger.DeduplicationStrategy.ONLY_MATCHING); |
||||
|
return trigger; |
||||
|
} |
||||
|
|
||||
|
private NotificationRule mockRule() { |
||||
|
NotificationRule rule = mock(NotificationRule.class); |
||||
|
when(rule.getDeduplicationKey()).thenReturn("rule:key"); |
||||
|
return rule; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,115 @@ |
|||||
|
/** |
||||
|
* 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.transport.lwm2m.bootstrap; |
||||
|
|
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.mockito.junit.jupiter.MockitoSettings; |
||||
|
import org.mockito.quality.Strictness; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.transport.TransportService; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.TbLwM2MDtlsBootstrapCertificateVerifier; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MBootstrapSecurityStore; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MInMemoryBootstrapConfigStore; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
||||
|
|
||||
|
import java.net.DatagramSocket; |
||||
|
import java.net.InetAddress; |
||||
|
import java.net.InetSocketAddress; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
@MockitoSettings(strictness = Strictness.LENIENT) |
||||
|
public class LwM2MTransportBootstrapServiceTest { |
||||
|
|
||||
|
private static final String HOST = "127.0.0.1"; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MTransportServerConfig serverConfig; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MTransportBootstrapConfig bootstrapConfig; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MBootstrapSecurityStore lwM2MBootstrapSecurityStore; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MInMemoryBootstrapConfigStore lwM2MInMemoryBootstrapConfigStore; |
||||
|
|
||||
|
@Mock |
||||
|
private TransportService transportService; |
||||
|
|
||||
|
@Mock |
||||
|
private TbLwM2MDtlsBootstrapCertificateVerifier certificateVerifier; |
||||
|
|
||||
|
private LwM2MTransportBootstrapService service; |
||||
|
private DatagramSocket occupiedPlain; |
||||
|
private DatagramSocket occupiedSecure; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() throws Exception { |
||||
|
occupiedPlain = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0)); |
||||
|
occupiedSecure = new DatagramSocket(new InetSocketAddress(InetAddress.getByName(HOST), 0)); |
||||
|
|
||||
|
when(bootstrapConfig.getHost()).thenReturn(HOST); |
||||
|
when(bootstrapConfig.getPort()).thenReturn(occupiedPlain.getLocalPort()); |
||||
|
when(bootstrapConfig.getSecureHost()).thenReturn(HOST); |
||||
|
when(bootstrapConfig.getSecurePort()).thenReturn(occupiedSecure.getLocalPort()); |
||||
|
when(bootstrapConfig.getSslCredentials()).thenReturn(null); |
||||
|
|
||||
|
when(serverConfig.isRecommendedCiphers()).thenReturn(false); |
||||
|
when(serverConfig.isRecommendedSupportedGroups()).thenReturn(false); |
||||
|
when(serverConfig.getDtlsRetransmissionTimeout()).thenReturn(9000); |
||||
|
when(serverConfig.getDtlsCidLength()).thenReturn(null); |
||||
|
|
||||
|
service = new LwM2MTransportBootstrapService( |
||||
|
serverConfig, |
||||
|
bootstrapConfig, |
||||
|
lwM2MBootstrapSecurityStore, |
||||
|
lwM2MInMemoryBootstrapConfigStore, |
||||
|
transportService, |
||||
|
certificateVerifier |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@AfterEach |
||||
|
public void tearDown() { |
||||
|
if (occupiedPlain != null && !occupiedPlain.isClosed()) { |
||||
|
occupiedPlain.close(); |
||||
|
} |
||||
|
if (occupiedSecure != null && !occupiedSecure.isClosed()) { |
||||
|
occupiedSecure.close(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void whenEndpointsFailToStart_thenInitThrowsAndReleasesBootstrapServer() { |
||||
|
assertThatThrownBy(() -> service.init()) |
||||
|
.isInstanceOf(IllegalStateException.class) |
||||
|
.hasMessageContaining("None of the server endpoints could be started"); |
||||
|
|
||||
|
assertThat(ReflectionTestUtils.getField(service, "server")).isNull(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,198 @@ |
|||||
|
/** |
||||
|
* 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.transport.lwm2m.bootstrap; |
||||
|
|
||||
|
import org.eclipse.leshan.server.bootstrap.LeshanBootstrapServer; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.Mockito; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.mockito.junit.jupiter.MockitoSettings; |
||||
|
import org.mockito.quality.Strictness; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.transport.TransportService; |
||||
|
import org.thingsboard.server.common.transport.config.ssl.SslCredentials; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.TbLwM2MDtlsBootstrapCertificateVerifier; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MBootstrapSecurityStore; |
||||
|
import org.thingsboard.server.transport.lwm2m.bootstrap.store.LwM2MInMemoryBootstrapConfigStore; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig; |
||||
|
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.doReturn; |
||||
|
import static org.mockito.Mockito.doThrow; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
@MockitoSettings(strictness = Strictness.LENIENT) |
||||
|
public class LwM2mBootstrapCertificateReloadTest { |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MTransportServerConfig mockServerConfig; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MTransportBootstrapConfig mockBootstrapConfig; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MBootstrapSecurityStore mockSecurityStore; |
||||
|
|
||||
|
@Mock |
||||
|
private LwM2MInMemoryBootstrapConfigStore mockConfigStore; |
||||
|
|
||||
|
@Mock |
||||
|
private TransportService mockTransportService; |
||||
|
|
||||
|
@Mock |
||||
|
private TbLwM2MDtlsBootstrapCertificateVerifier mockCertificateVerifier; |
||||
|
|
||||
|
@Mock |
||||
|
private LeshanBootstrapServer mockBootstrapServer; |
||||
|
|
||||
|
@Mock |
||||
|
private SslCredentials mockSslCredentials; |
||||
|
|
||||
|
private LwM2MTransportBootstrapService bootstrapService; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setup() { |
||||
|
bootstrapService = new LwM2MTransportBootstrapService( |
||||
|
mockServerConfig, |
||||
|
mockBootstrapConfig, |
||||
|
mockSecurityStore, |
||||
|
mockConfigStore, |
||||
|
mockTransportService, |
||||
|
mockCertificateVerifier |
||||
|
); |
||||
|
|
||||
|
when(mockBootstrapConfig.getHost()).thenReturn("localhost"); |
||||
|
when(mockBootstrapConfig.getPort()).thenReturn(5687); |
||||
|
when(mockBootstrapConfig.getSecureHost()).thenReturn("localhost"); |
||||
|
when(mockBootstrapConfig.getSecurePort()).thenReturn(5688); |
||||
|
when(mockBootstrapConfig.getSslCredentials()).thenReturn(mockSslCredentials); |
||||
|
when(mockServerConfig.getDtlsRetransmissionTimeout()).thenReturn(9000); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenInit_whenCalled_thenShouldRegisterCertificateReloadCallback() { |
||||
|
ReflectionTestUtils.setField(bootstrapService, "server", mockBootstrapServer); |
||||
|
|
||||
|
bootstrapService.afterSingletonsInstantiated(); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
verify(mockBootstrapConfig).registerServerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
assertThat(callbackCaptor.getValue()).isNotNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenNewServerCreationFails_thenOldServerIsPreserved() { |
||||
|
ReflectionTestUtils.setField(bootstrapService, "server", mockBootstrapServer); |
||||
|
|
||||
|
// Force getLhBootstrapServer() to fail by returning null host (causes InetSocketAddress to throw)
|
||||
|
when(mockBootstrapConfig.getHost()).thenReturn(null); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
bootstrapService.afterSingletonsInstantiated(); |
||||
|
verify(mockBootstrapConfig).registerServerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
|
||||
|
// getLhBootstrapServer() will fail due to null host before old server is stopped.
|
||||
|
// The old server should NOT be destroyed since the new server was never created.
|
||||
|
reloadCallback.run(); |
||||
|
|
||||
|
verify(mockBootstrapServer, never()).stop(); |
||||
|
verify(mockBootstrapServer, never()).destroy(); |
||||
|
assertThat(ReflectionTestUtils.getField(bootstrapService, "server")).isSameAs(mockBootstrapServer); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenNullServer_whenRecreate_thenShouldNotThrow() { |
||||
|
ReflectionTestUtils.setField(bootstrapService, "server", null); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
bootstrapService.afterSingletonsInstantiated(); |
||||
|
verify(mockBootstrapConfig).registerServerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
|
||||
|
// Should not throw — callback catches exceptions internally
|
||||
|
reloadCallback.run(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenCertificateUpdate_whenRecreate_thenShouldUseNewCredentials() { |
||||
|
SslCredentials oldCredentials = mockSslCredentials; |
||||
|
SslCredentials newCredentials = mock(SslCredentials.class); |
||||
|
|
||||
|
when(mockBootstrapConfig.getSslCredentials()).thenReturn(oldCredentials).thenReturn(newCredentials); |
||||
|
|
||||
|
SslCredentials firstCall = mockBootstrapConfig.getSslCredentials(); |
||||
|
assertThat(firstCall).isEqualTo(oldCredentials); |
||||
|
|
||||
|
SslCredentials secondCall = mockBootstrapConfig.getSslCredentials(); |
||||
|
assertThat(secondCall).isEqualTo(newCredentials); |
||||
|
|
||||
|
verify(mockBootstrapConfig, times(2)).getSslCredentials(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenRegistered_thenShouldRegisterExactlyOne() { |
||||
|
bootstrapService.afterSingletonsInstantiated(); |
||||
|
|
||||
|
verify(mockBootstrapConfig, times(1)).registerServerReloadCallback(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void givenReloadCallback_whenNewServerStartFails_thenOldServerRestarted() { |
||||
|
// GIVEN
|
||||
|
ReflectionTestUtils.setField(bootstrapService, "server", mockBootstrapServer); |
||||
|
|
||||
|
LeshanBootstrapServer mockNewServer = mock(LeshanBootstrapServer.class); |
||||
|
doThrow(new RuntimeException("start failed")).when(mockNewServer).start(); |
||||
|
|
||||
|
LwM2MTransportBootstrapService spyService = Mockito.spy(bootstrapService); |
||||
|
doReturn(mockNewServer).when(spyService).getLhBootstrapServer(); |
||||
|
|
||||
|
ArgumentCaptor<Runnable> callbackCaptor = ArgumentCaptor.forClass(Runnable.class); |
||||
|
spyService.afterSingletonsInstantiated(); |
||||
|
verify(mockBootstrapConfig).registerServerReloadCallback(callbackCaptor.capture()); |
||||
|
|
||||
|
Runnable reloadCallback = callbackCaptor.getValue(); |
||||
|
|
||||
|
// WHEN
|
||||
|
reloadCallback.run(); |
||||
|
|
||||
|
// THEN
|
||||
|
// Old server is stopped (not destroyed) to release ports
|
||||
|
verify(mockBootstrapServer).stop(); |
||||
|
verify(mockBootstrapServer, never()).destroy(); |
||||
|
// The new server fails to start and is destroyed
|
||||
|
verify(mockNewServer).destroy(); |
||||
|
// Old server is restarted (not rebuilt from potentially stale credentials)
|
||||
|
verify(mockBootstrapServer).start(); |
||||
|
assertThat(ReflectionTestUtils.getField(spyService, "server")).isSameAs(mockBootstrapServer); |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue