Browse Source
# Conflicts: # application/pom.xml # common/actor/pom.xml # common/cache/pom.xml # common/cluster-api/pom.xml # common/coap-server/pom.xml # common/dao-api/pom.xml # common/data/pom.xml # common/discovery-api/pom.xml # common/edge-api/pom.xml # common/edqs/pom.xml # common/message/pom.xml # common/pom.xml # common/proto/pom.xml # common/queue/pom.xml # common/script/pom.xml # common/script/remote-js-client/pom.xml # common/script/script-api/pom.xml # common/stats/pom.xml # common/transport/coap/pom.xml # common/transport/http/pom.xml # common/transport/lwm2m/pom.xml # common/transport/mqtt/pom.xml # common/transport/pom.xml # common/transport/snmp/pom.xml # common/transport/transport-api/pom.xml # common/util/pom.xml # common/version-control/pom.xml # dao/pom.xml # edqs/pom.xml # monitoring/pom.xml # msa/black-box-tests/pom.xml # msa/edqs/pom.xml # msa/js-executor/pom.xml # msa/monitoring/pom.xml # msa/pom.xml # msa/tb-node/pom.xml # msa/tb/pom.xml # msa/transport/coap/pom.xml # msa/transport/http/pom.xml # msa/transport/lwm2m/pom.xml # msa/transport/mqtt/pom.xml # msa/transport/pom.xml # msa/transport/snmp/pom.xml # msa/vc-executor-docker/pom.xml # msa/vc-executor/pom.xml # msa/web-ui/pom.xml # netty-mqtt/pom.xml # pom.xml # rest-client/pom.xml # rule-engine/pom.xml # rule-engine/rule-engine-api/pom.xml # rule-engine/rule-engine-components/pom.xml # tools/pom.xml # transport/coap/pom.xml # transport/http/pom.xml # transport/lwm2m/pom.xml # transport/mqtt/pom.xml # transport/pom.xml # transport/snmp/pom.xml # ui-ngx/pom.xmlrelease-4.2
263 changed files with 12066 additions and 2098 deletions
File diff suppressed because one or more lines are too long
@ -0,0 +1,19 @@ |
|||||
|
-- |
||||
|
-- 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. |
||||
|
-- |
||||
|
|
||||
|
-- LTS cumulative schema update file. |
||||
|
-- All statements must be idempotent (use IF NOT EXISTS, ADD COLUMN IF NOT EXISTS, DO $$ ... END $$ guards, etc.). |
||||
|
-- This file is executed by SystemPatchApplier on every version increase within the LTS family. |
||||
@ -0,0 +1,51 @@ |
|||||
|
/** |
||||
|
* 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.config; |
||||
|
|
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.rule.engine.api.TbHttpClientSettings; |
||||
|
import org.thingsboard.server.queue.util.TbRuleEngineComponent; |
||||
|
|
||||
|
@TbRuleEngineComponent |
||||
|
@Component |
||||
|
public class TbHttpClientSettingsComponent implements TbHttpClientSettings { |
||||
|
|
||||
|
@Value("${actors.rule.external.http_client.max_parallel_requests:0}") |
||||
|
private int maxParallelRequests; |
||||
|
|
||||
|
@Value("${actors.rule.external.http_client.max_pending_requests:0}") |
||||
|
private int maxPendingRequests; |
||||
|
|
||||
|
@Value("${actors.rule.external.http_client.pool_max_connections:0}") |
||||
|
private int poolMaxConnections; |
||||
|
|
||||
|
@Override |
||||
|
public int getMaxParallelRequests() { |
||||
|
return maxParallelRequests; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public int getMaxPendingRequests() { |
||||
|
return maxPendingRequests; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public int getPoolMaxConnections() { |
||||
|
return poolMaxConnections; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,144 @@ |
|||||
|
/** |
||||
|
* 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.ai; |
||||
|
|
||||
|
import com.google.cloud.vertexai.api.GenerationConfig; |
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.parallel.ResourceLock; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.common.util.SsrfProtectionValidator; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AzureOpenAiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.GoogleVertexAiGeminiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.OllamaChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.OpenAiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.AzureOpenAiProviderConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.GoogleVertexAiGeminiProviderConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.OllamaProviderConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.OpenAiProviderConfig; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
|
||||
|
@ResourceLock("SsrfProtectionValidator") |
||||
|
class Langchain4jChatModelConfigurerImplTest { |
||||
|
|
||||
|
private static final String TEST_SERVICE_ACCOUNT_KEY = """ |
||||
|
{ |
||||
|
"type": "service_account", |
||||
|
"project_id": "test-project", |
||||
|
"private_key_id": "key-id", |
||||
|
"private_key": "-----BEGIN PRIVATE KEY-----\\nMIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDNrHph/y7zyxIg\\ncmYYeOD8mFg9KraK71n84ffTQyVrl4HzlQgRIz5m4vM2rV5zjVLFi0xlAPT/iq/5\\nbh2zA4iXI0dEsR901nVjcL182t/GRYbKen53ZiuScBxBoCZPXW16Md+Yk8nMdNUb\\n4LoIRGZq64bjsJ+vh3Aa2gdGUHpDyebIXlXbF8ehWmEhgUsL7XjL0PkJ4lt2UMG+\\nx1j2Or25rqJmfc5M9kbxvINtdvSRTPiMOIXX00fCDZjQdd18RBVHOxraGxDgQpmv\\nk4qjFEPqGr0YTsa5dI8fz+4DqJpEi3rancRiTKM/KUwYLGnPSD07XqGfiDA8npVm\\nj4N62LhnAgMBAAECggEADbFfH87DRk7YQO8XgKdCOf7oglX+0NwjjmmlkVvwjgEI\\nZqT0ObPcz9u/MSjfV2vAEs/LK773ELD1NfLQqQjiBfHpfkIZOTLynhwKOYRBjqvf\\n+p38ynUzucGbV/vSC8meuW/AQPe3Nn9MFYQ4znEYrSNLbTWRRA3idvSEtHfffqDA\\nDHRBI1eIlxh1OTIR3L+HhcNYuus1LuoKnSlmwLGhAZLt7fjuWK3PkOiFT15e0M9M\\nUhp3WwhHcRC0o6bxT+BWRYKMVX3Vjlro4sF1fq4+jePThX1bpJcPfUmsC95tXPfX\\njfNGAxHlZ+MS1V/cLlIqyz7drXBcwCDJtbPmvavmNQKBgQD1ZR/ePHcjXUGM57U8\\nbxPatNOrcicvaP2AtTA6Y/JjfbcydkVXsenDGk0h6hykpIiMwrAaJuUTjGeM6QTI\\nOhK0k1QbGitcM71d9TSdLzWUdb3yvsyaPZlPR/6u2FBb6Bf9rWOQRYYyv1Lvu0+1\\nYLnR4sHBxiAur1NGxuHfA4ZUOwKBgQDWj+fcS/x/ifbCqexr7teWU+tUeyG3eLGA\\nMmB9eCkY7djl0/LHu/IHgqrGRVgra3IB1uI7Wr3jZYvlS7qGL3KpjeIPYj7LTQC5\\nznm0875NvJELPjQK/A4EM3mC057QRvb7y52KBNKJi7/JwHU7VHmudB78e7uGlW2K\\n5Ccl0PJFxQKBgCWv5yoJXT64JsYOG95xLLptBQkSmgQE+tHWgdal3Ob8urLsSRAD\\nyePl2Sy5OLbscfA0Qjlx+cJ70LdqXgqmKJNFASi8ZyZc59tTOkZdprvrLUXnmaKi\\njTYI14tgu06yIWUbSOwyUT7f9UvOF5rChSc/zQQGepDQ6lg3WR8X+nxbAoGBAKiu\\nfAcqSfjuuuuxcWgtXpoVoaZKI2i9Xza85DTf+ddabjHJXk3+iTm0VZQIwldoYjnl\\n+PfW0ABtPf1net2xgcChBf84Ksvj3tU06WQEWDF/NLyVC48zN8W/viDHREzT7app\\nGpJ+VhLCpmXzg3bAY+Vt70pp8DTPV05hLhHB4iZNAoGBAJW+bYh7jE61+58VpjvF\\nP36BK09jEEPWVucJdghb2mb62iA6JDy3ApU+8FzckXHDewt0sqvsW4VqukgwVZx3\\npSC7mR4B+Fm6znm0Z5mBWiG5bOOgTJ0mRZv4cYgC+JRRF/E3yYR58RyAKFAFIAFH\\nng8XYP1wQp64Fzv4+rUSwM49\\n-----END PRIVATE KEY-----\\n", |
||||
|
"client_email": "test@test-project.iam.gserviceaccount.com", |
||||
|
"client_id": "123456789", |
||||
|
"auth_uri": "https://accounts.google.com/o/oauth2/auth", |
||||
|
"token_uri": "https://oauth2.googleapis.com/token" |
||||
|
} |
||||
|
"""; |
||||
|
|
||||
|
private final Langchain4jChatModelConfigurerImpl configurer = new Langchain4jChatModelConfigurerImpl(); |
||||
|
|
||||
|
@BeforeEach |
||||
|
void enableSsrfProtection() { |
||||
|
SsrfProtectionValidator.setEnabled(true); |
||||
|
} |
||||
|
|
||||
|
@AfterEach |
||||
|
void disableSsrfProtection() { |
||||
|
SsrfProtectionValidator.setEnabled(false); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void configureChatModel_openAi_withPrivateIp_shouldThrow() { |
||||
|
var config = OpenAiChatModelConfig.builder() |
||||
|
.providerConfig(OpenAiProviderConfig.builder() |
||||
|
.baseUrl("http://172.17.0.1:8080/") |
||||
|
.apiKey("test") |
||||
|
.build()) |
||||
|
.modelId("gpt-4o") |
||||
|
.build(); |
||||
|
|
||||
|
assertThatThrownBy(() -> configurer.configureChatModel(config)) |
||||
|
.isInstanceOf(RuntimeException.class) |
||||
|
.hasMessageContaining("URI is invalid"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void configureChatModel_openAi_withLocalhostUrl_shouldThrow() { |
||||
|
var config = OpenAiChatModelConfig.builder() |
||||
|
.providerConfig(OpenAiProviderConfig.builder() |
||||
|
.baseUrl("http://localhost:22/") |
||||
|
.apiKey("test") |
||||
|
.build()) |
||||
|
.modelId("gpt-4o") |
||||
|
.build(); |
||||
|
|
||||
|
assertThatThrownBy(() -> configurer.configureChatModel(config)) |
||||
|
.isInstanceOf(RuntimeException.class) |
||||
|
.hasMessageContaining("URI is invalid"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void configureChatModel_azureOpenAi_withPrivateIp_shouldThrow() { |
||||
|
var config = AzureOpenAiChatModelConfig.builder() |
||||
|
.providerConfig(new AzureOpenAiProviderConfig( |
||||
|
"http://10.0.0.1:8080/", null, "test-key")) |
||||
|
.modelId("gpt-4o") |
||||
|
.build(); |
||||
|
|
||||
|
assertThatThrownBy(() -> configurer.configureChatModel(config)) |
||||
|
.isInstanceOf(RuntimeException.class) |
||||
|
.hasMessageContaining("URI is invalid"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void configureChatModel_ollama_withPrivateIp_shouldThrow() { |
||||
|
var config = OllamaChatModelConfig.builder() |
||||
|
.providerConfig(new OllamaProviderConfig( |
||||
|
"http://192.168.1.100:11434/", new OllamaProviderConfig.OllamaAuth.None())) |
||||
|
.modelId("llama3") |
||||
|
.build(); |
||||
|
|
||||
|
assertThatThrownBy(() -> configurer.configureChatModel(config)) |
||||
|
.isInstanceOf(RuntimeException.class) |
||||
|
.hasMessageContaining("URI is invalid"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void configureChatModel_vertexAi_setsFrequencyAndPresencePenaltyFromCorrectConfigFields() { |
||||
|
// GIVEN
|
||||
|
var providerConfig = new GoogleVertexAiGeminiProviderConfig( |
||||
|
"test.json", "test-project", "us-central1", TEST_SERVICE_ACCOUNT_KEY |
||||
|
); |
||||
|
var chatModelConfig = GoogleVertexAiGeminiChatModelConfig.builder() |
||||
|
.providerConfig(providerConfig) |
||||
|
.modelId("gemini-2.0-flash") |
||||
|
.frequencyPenalty(0.3) |
||||
|
.presencePenalty(0.7) |
||||
|
.build(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ChatModel chatModel = configurer.configureChatModel(chatModelConfig); |
||||
|
|
||||
|
// THEN
|
||||
|
var generationConfig = (GenerationConfig) ReflectionTestUtils.getField(chatModel, "generationConfig"); |
||||
|
assertThat(generationConfig.getFrequencyPenalty()).isEqualTo(0.3f); |
||||
|
assertThat(generationConfig.getPresencePenalty()).isEqualTo(0.7f); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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.service.cf; |
||||
|
|
||||
|
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.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedFieldLink; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.id.AssetProfileId; |
||||
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
||||
|
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
||||
|
import org.thingsboard.server.dao.asset.AssetService; |
||||
|
import org.thingsboard.server.dao.cf.CalculatedFieldService; |
||||
|
import org.thingsboard.server.dao.customer.CustomerService; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
|
||||
|
import java.util.Collections; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultCalculatedFieldCacheTest { |
||||
|
|
||||
|
@Mock |
||||
|
private CalculatedFieldService calculatedFieldService; |
||||
|
@Mock |
||||
|
private DeviceService deviceService; |
||||
|
@Mock |
||||
|
private AssetService assetService; |
||||
|
@Mock |
||||
|
private CustomerService customerService; |
||||
|
|
||||
|
private DefaultCalculatedFieldCache cache; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
cache = new DefaultCalculatedFieldCache(calculatedFieldService, null, null); |
||||
|
} |
||||
|
|
||||
|
// --- Tenant deletion tests ---
|
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_evictsAllTenantCfsFromAllMaps() { |
||||
|
TenantId tenant1 = new TenantId(UUID.randomUUID()); |
||||
|
TenantId tenant2 = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device1 = new DeviceId(UUID.randomUUID()); |
||||
|
DeviceId device2 = new DeviceId(UUID.randomUUID()); |
||||
|
|
||||
|
CalculatedField cf1 = addCfToCache(tenant1, device1); |
||||
|
CalculatedField cf2 = addCfToCache(tenant2, device2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant1, tenant1, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf1.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(device1)).isEmpty(); |
||||
|
assertThat(cache.getCalculatedField(cf2.getId())).isEqualTo(cf2); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(device2)).containsExactly(cf2); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_removesLinksToLinkedEntities() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId cfEntity = new DeviceId(UUID.randomUUID()); |
||||
|
DeviceId linkedDevice = new DeviceId(UUID.randomUUID()); |
||||
|
|
||||
|
CalculatedField cf = addCfToCache(tenant, cfEntity, linkedDevice); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedFieldLinksByEntityId(linkedDevice)).isEmpty(); |
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantUpdated_doesNotEvictCfs() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, device); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isEqualTo(cf); |
||||
|
} |
||||
|
|
||||
|
// --- Device/Asset deletion tests ---
|
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceDeleted_evictsCfsForThatDevice() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, device); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, device, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(device)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceDeleted_removesLinksForLinkedEntities() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
DeviceId linkedDevice = new DeviceId(UUID.randomUUID()); |
||||
|
addCfToCache(tenant, device, linkedDevice); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, device, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedFieldLinksByEntityId(linkedDevice)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_assetDeleted_evictsCfsForThatAsset() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetId asset = new AssetId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, asset); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, asset, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(asset)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
// --- DeviceProfile/AssetProfile deletion tests ---
|
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceProfileDeleted_evictsCfsForThatProfile() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profileId)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceProfileDeleted_removesLinksForLinkedEntities() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
DeviceId linkedDevice = new DeviceId(UUID.randomUUID()); |
||||
|
addCfToCache(tenant, profileId, linkedDevice); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedFieldLinksByEntityId(linkedDevice)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceProfileDeleted_doesNotEvictOtherProfilesCfs() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profile1 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
DeviceProfileId profile2 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf1 = addCfToCache(tenant, profile1); |
||||
|
CalculatedField cf2 = addCfToCache(tenant, profile2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profile1, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf1.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profile1)).isEmpty(); |
||||
|
assertThat(cache.getCalculatedField(cf2.getId())).isEqualTo(cf2); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profile2)).containsExactly(cf2); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_deviceProfileUpdated_doesNotEvictCfs() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isEqualTo(cf); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profileId)).containsExactly(cf); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_assetProfileDeleted_evictsCfsForThatProfile() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profileId)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_assetProfileDeleted_removesLinksForLinkedEntities() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
AssetId linkedAsset = new AssetId(UUID.randomUUID()); |
||||
|
addCfToCache(tenant, profileId, linkedAsset); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedFieldLinksByEntityId(linkedAsset)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_assetProfileDeleted_doesNotEvictOtherProfilesCfs() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profile1 = new AssetProfileId(UUID.randomUUID()); |
||||
|
AssetProfileId profile2 = new AssetProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf1 = addCfToCache(tenant, profile1); |
||||
|
CalculatedField cf2 = addCfToCache(tenant, profile2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profile1, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf1.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profile1)).isEmpty(); |
||||
|
assertThat(cache.getCalculatedField(cf2.getId())).isEqualTo(cf2); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profile2)).containsExactly(cf2); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_assetProfileUpdated_doesNotEvictCfs() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, profileId, ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isEqualTo(cf); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(profileId)).containsExactly(cf); |
||||
|
} |
||||
|
|
||||
|
// --- CalculatedField lifecycle tests ---
|
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_calculatedFieldCreated_addsCfToCache() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
CalculatedFieldId cfId = new CalculatedFieldId(UUID.randomUUID()); |
||||
|
CalculatedField cf = buildCalculatedField(cfId, tenant, device, simpleCfConfig()); |
||||
|
when(calculatedFieldService.findById(tenant, cfId)).thenReturn(cf); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, cfId, ComponentLifecycleEvent.CREATED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cfId)).isEqualTo(cf); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(device)).containsExactly(cf); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_calculatedFieldDeleted_evictsCfFromCache() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, device); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, cf.getId(), ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isNull(); |
||||
|
assertThat(cache.getCalculatedFieldsByEntityId(device)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_calculatedFieldUpdated_refreshesCfInCache() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId device = new DeviceId(UUID.randomUUID()); |
||||
|
CalculatedField cf = addCfToCache(tenant, device); |
||||
|
|
||||
|
CalculatedField updatedCf = buildCalculatedField(cf.getId(), tenant, device, simpleCfConfig()); |
||||
|
updatedCf.setName("updated-name"); |
||||
|
when(calculatedFieldService.findById(tenant, cf.getId())).thenReturn(updatedCf); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, cf.getId(), ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
assertThat(cache.getCalculatedField(cf.getId())).isEqualTo(updatedCf); |
||||
|
} |
||||
|
|
||||
|
private CalculatedField addCfToCache(TenantId tenantId, EntityId entityId) { |
||||
|
CalculatedFieldId cfId = new CalculatedFieldId(UUID.randomUUID()); |
||||
|
CalculatedField cf = buildCalculatedField(cfId, tenantId, entityId, simpleCfConfig()); |
||||
|
when(calculatedFieldService.findById(tenantId, cfId)).thenReturn(cf); |
||||
|
cache.addCalculatedField(tenantId, cfId); |
||||
|
return cf; |
||||
|
} |
||||
|
|
||||
|
private CalculatedField addCfToCache(TenantId tenantId, EntityId entityId, EntityId linkedEntity) { |
||||
|
CalculatedFieldId cfId = new CalculatedFieldId(UUID.randomUUID()); |
||||
|
CalculatedFieldConfiguration config = linkedEntityCfConfig(tenantId, cfId, linkedEntity); |
||||
|
CalculatedField cf = buildCalculatedField(cfId, tenantId, entityId, config); |
||||
|
when(calculatedFieldService.findById(tenantId, cfId)).thenReturn(cf); |
||||
|
cache.addCalculatedField(tenantId, cfId); |
||||
|
return cf; |
||||
|
} |
||||
|
|
||||
|
private CalculatedField buildCalculatedField(CalculatedFieldId id, TenantId tenantId, EntityId entityId, CalculatedFieldConfiguration config) { |
||||
|
CalculatedField cf = new CalculatedField(); |
||||
|
cf.setId(id); |
||||
|
cf.setTenantId(tenantId); |
||||
|
cf.setEntityId(entityId); |
||||
|
cf.setType(CalculatedFieldType.SIMPLE); |
||||
|
cf.setName("test-cf-" + id.getId()); |
||||
|
cf.setConfiguration(config); |
||||
|
return cf; |
||||
|
} |
||||
|
|
||||
|
private CalculatedFieldConfiguration simpleCfConfig() { |
||||
|
CalculatedFieldConfiguration config = mock(CalculatedFieldConfiguration.class); |
||||
|
when(config.getReferencedEntities()).thenReturn(Collections.emptyList()); |
||||
|
when(config.buildCalculatedFieldLinks(any(), any(), any())).thenReturn(Collections.emptyList()); |
||||
|
return config; |
||||
|
} |
||||
|
|
||||
|
private CalculatedFieldConfiguration linkedEntityCfConfig(TenantId tenantId, CalculatedFieldId cfId, EntityId linkedEntity) { |
||||
|
CalculatedFieldConfiguration config = mock(CalculatedFieldConfiguration.class); |
||||
|
CalculatedFieldLink link = new CalculatedFieldLink(tenantId, linkedEntity, cfId); |
||||
|
when(config.getReferencedEntities()).thenReturn(List.of(linkedEntity)); |
||||
|
when(config.buildCalculatedFieldLinks(any(), any(), any())).thenReturn(List.of(link)); |
||||
|
when(config.buildCalculatedFieldLink(any(), eq(linkedEntity), any())).thenReturn(link); |
||||
|
return config; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,159 @@ |
|||||
|
/** |
||||
|
* 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.profile; |
||||
|
|
||||
|
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.thingsboard.server.common.data.asset.Asset; |
||||
|
import org.thingsboard.server.common.data.asset.AssetProfile; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.id.AssetProfileId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
||||
|
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
||||
|
import org.thingsboard.server.dao.asset.AssetProfileService; |
||||
|
import org.thingsboard.server.dao.asset.AssetService; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultTbAssetProfileCacheTest { |
||||
|
|
||||
|
@Mock |
||||
|
private AssetProfileService assetProfileService; |
||||
|
@Mock |
||||
|
private AssetService assetService; |
||||
|
|
||||
|
private DefaultTbAssetProfileCache cache; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
cache = new DefaultTbAssetProfileCache(assetProfileService, assetService); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_evictsAssetProfilesForThatTenant() { |
||||
|
TenantId tenant1 = new TenantId(UUID.randomUUID()); |
||||
|
TenantId tenant2 = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId1 = new AssetProfileId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId2 = new AssetProfileId(UUID.randomUUID()); |
||||
|
|
||||
|
loadProfileIntoCache(tenant1, profileId1); |
||||
|
loadProfileIntoCache(tenant2, profileId2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant1, tenant1, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// After deletion tenant1 profile should be reloaded from service on next get
|
||||
|
when(assetProfileService.findAssetProfileById(any(), any())).thenReturn(null); |
||||
|
assertThat(cache.get(tenant1, profileId1)).isNull(); |
||||
|
verify(assetProfileService, times(1)).findAssetProfileById(tenant2, profileId2); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_evictsAssetMappingsForThatTenant() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
AssetId assetId = new AssetId(UUID.randomUUID()); |
||||
|
|
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
loadAssetMappingIntoCache(tenant, assetId, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// After tenant deletion, asset-to-profile mapping should be gone; get() should try to reload
|
||||
|
when(assetService.findAssetById(any(), any())).thenReturn(null); |
||||
|
assertThat(cache.get(tenant, assetId)).isNull(); |
||||
|
verify(assetService, times(2)).findAssetById(tenant, assetId); // once on load, once after eviction
|
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_removesListenersForThatTenant() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
EntityId listenerId = new AssetId(UUID.randomUUID()); |
||||
|
AtomicInteger callCount = new AtomicInteger(); |
||||
|
|
||||
|
cache.addListener(tenant, listenerId, profile -> callCount.incrementAndGet(), null); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// Evicting a profile after tenant deletion should not trigger the removed listener
|
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
cache.evict(tenant, profileId); |
||||
|
|
||||
|
assertThat(callCount.get()).isZero(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantUpdated_doesNotEvictProfiles() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId = new AssetProfileId(UUID.randomUUID()); |
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
// Profile should still be served from cache without hitting the service again
|
||||
|
cache.get(tenant, profileId); |
||||
|
verify(assetProfileService, times(1)).findAssetProfileById(tenant, profileId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_differentTenantDeleted_keepsOtherTenantsProfiles() { |
||||
|
TenantId tenant1 = new TenantId(UUID.randomUUID()); |
||||
|
TenantId tenant2 = new TenantId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId1 = new AssetProfileId(UUID.randomUUID()); |
||||
|
AssetProfileId profileId2 = new AssetProfileId(UUID.randomUUID()); |
||||
|
|
||||
|
AssetProfile profile1 = loadProfileIntoCache(tenant1, profileId1); |
||||
|
loadProfileIntoCache(tenant2, profileId2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant2, tenant2, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.get(tenant1, profileId1)).isEqualTo(profile1); |
||||
|
verify(assetProfileService, times(1)).findAssetProfileById(tenant1, profileId1); |
||||
|
} |
||||
|
|
||||
|
// --- Helpers ---
|
||||
|
|
||||
|
private AssetProfile loadProfileIntoCache(TenantId tenantId, AssetProfileId profileId) { |
||||
|
AssetProfile profile = new AssetProfile(); |
||||
|
profile.setId(profileId); |
||||
|
profile.setTenantId(tenantId); |
||||
|
when(assetProfileService.findAssetProfileById(tenantId, profileId)).thenReturn(profile); |
||||
|
cache.get(tenantId, profileId); |
||||
|
return profile; |
||||
|
} |
||||
|
|
||||
|
private void loadAssetMappingIntoCache(TenantId tenantId, AssetId assetId, AssetProfileId profileId) { |
||||
|
Asset asset = new Asset(); |
||||
|
asset.setId(assetId); |
||||
|
asset.setAssetProfileId(profileId); |
||||
|
when(assetService.findAssetById(tenantId, assetId)).thenReturn(asset); |
||||
|
cache.get(tenantId, assetId); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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.service.profile; |
||||
|
|
||||
|
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.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.DeviceProfile; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
||||
|
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
||||
|
import org.thingsboard.server.dao.device.DeviceProfileService; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class DefaultTbDeviceProfileCacheTest { |
||||
|
|
||||
|
@Mock |
||||
|
private DeviceProfileService deviceProfileService; |
||||
|
@Mock |
||||
|
private DeviceService deviceService; |
||||
|
|
||||
|
private DefaultTbDeviceProfileCache cache; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
cache = new DefaultTbDeviceProfileCache(deviceProfileService, deviceService); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_evictsDeviceProfilesForThatTenant() { |
||||
|
TenantId tenant1 = new TenantId(UUID.randomUUID()); |
||||
|
TenantId tenant2 = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId1 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId2 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
|
||||
|
loadProfileIntoCache(tenant1, profileId1); |
||||
|
loadProfileIntoCache(tenant2, profileId2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant1, tenant1, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// After deletion tenant1 profile should be reloaded from service on next get
|
||||
|
when(deviceProfileService.findDeviceProfileById(any(), any())).thenReturn(null); |
||||
|
assertThat(cache.get(tenant1, profileId1)).isNull(); |
||||
|
// tenant2 profile should still be served from cache (no extra service call)
|
||||
|
verify(deviceProfileService, times(1)).findDeviceProfileById(tenant2, profileId2); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_evictsDeviceMappingsForThatTenant() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
||||
|
|
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
loadDeviceMappingIntoCache(tenant, deviceId, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// After tenant deletion, device-to-profile mapping should be gone; get() should try to reload
|
||||
|
when(deviceService.findDeviceById(any(), any())).thenReturn(null); |
||||
|
assertThat(cache.get(tenant, deviceId)).isNull(); |
||||
|
verify(deviceService, times(2)).findDeviceById(tenant, deviceId); // once on load, once after eviction
|
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantDeleted_removesListenersForThatTenant() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
EntityId listenerId = new DeviceId(UUID.randomUUID()); |
||||
|
AtomicInteger callCount = new AtomicInteger(); |
||||
|
|
||||
|
cache.addListener(tenant, listenerId, profile -> callCount.incrementAndGet(), null); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
// Evicting a profile after tenant deletion should not trigger the removed listener
|
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
cache.evict(tenant, profileId); |
||||
|
|
||||
|
assertThat(callCount.get()).isZero(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_tenantUpdated_doesNotEvictProfiles() { |
||||
|
TenantId tenant = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId = new DeviceProfileId(UUID.randomUUID()); |
||||
|
loadProfileIntoCache(tenant, profileId); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant, tenant, ComponentLifecycleEvent.UPDATED)); |
||||
|
|
||||
|
// Profile should still be served from cache without hitting the service again
|
||||
|
cache.get(tenant, profileId); |
||||
|
verify(deviceProfileService, times(1)).findDeviceProfileById(tenant, profileId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void onComponentLifecycleEvent_differentTenantDeleted_keepsOtherTenantsProfiles() { |
||||
|
TenantId tenant1 = new TenantId(UUID.randomUUID()); |
||||
|
TenantId tenant2 = new TenantId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId1 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
DeviceProfileId profileId2 = new DeviceProfileId(UUID.randomUUID()); |
||||
|
|
||||
|
DeviceProfile profile1 = loadProfileIntoCache(tenant1, profileId1); |
||||
|
loadProfileIntoCache(tenant2, profileId2); |
||||
|
|
||||
|
cache.onComponentLifecycleEvent(new ComponentLifecycleMsg(tenant2, tenant2, ComponentLifecycleEvent.DELETED)); |
||||
|
|
||||
|
assertThat(cache.get(tenant1, profileId1)).isEqualTo(profile1); |
||||
|
verify(deviceProfileService, times(1)).findDeviceProfileById(tenant1, profileId1); |
||||
|
} |
||||
|
|
||||
|
// --- Helpers ---
|
||||
|
|
||||
|
private DeviceProfile loadProfileIntoCache(TenantId tenantId, DeviceProfileId profileId) { |
||||
|
DeviceProfile profile = new DeviceProfile(); |
||||
|
profile.setId(profileId); |
||||
|
profile.setTenantId(tenantId); |
||||
|
when(deviceProfileService.findDeviceProfileById(tenantId, profileId)).thenReturn(profile); |
||||
|
cache.get(tenantId, profileId); |
||||
|
return profile; |
||||
|
} |
||||
|
|
||||
|
private void loadDeviceMappingIntoCache(TenantId tenantId, DeviceId deviceId, DeviceProfileId profileId) { |
||||
|
Device device = new Device(); |
||||
|
device.setId(deviceId); |
||||
|
device.setDeviceProfileId(profileId); |
||||
|
when(deviceService.findDeviceById(tenantId, deviceId)).thenReturn(device); |
||||
|
cache.get(tenantId, deviceId); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,277 @@ |
|||||
|
/** |
||||
|
* 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.ws; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.data.TenantProfile; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.id.UserId; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; |
||||
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
||||
|
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; |
||||
|
import org.thingsboard.server.service.security.AccessValidator; |
||||
|
import org.thingsboard.server.service.security.model.SecurityUser; |
||||
|
import org.thingsboard.server.service.security.model.UserPrincipal; |
||||
|
import org.thingsboard.server.service.subscription.TbEntityDataSubscriptionService; |
||||
|
import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; |
||||
|
import org.thingsboard.server.service.ws.notification.NotificationCommandsHandler; |
||||
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.AttributesSubscriptionCmd; |
||||
|
|
||||
|
import java.util.Set; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.BDDMockito.willReturn; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
|
||||
|
class DefaultWebSocketServiceTest { |
||||
|
|
||||
|
DefaultWebSocketService service; |
||||
|
TbTenantProfileCache tenantProfileCache; |
||||
|
WebSocketMsgEndpoint msgEndpoint; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
tenantProfileCache = mock(TbTenantProfileCache.class); |
||||
|
msgEndpoint = mock(WebSocketMsgEndpoint.class); |
||||
|
|
||||
|
service = new DefaultWebSocketService( |
||||
|
mock(TbLocalSubscriptionService.class), |
||||
|
mock(TbEntityDataSubscriptionService.class), |
||||
|
mock(NotificationCommandsHandler.class), |
||||
|
msgEndpoint, |
||||
|
mock(AccessValidator.class), |
||||
|
mock(AttributesService.class), |
||||
|
mock(TimeseriesService.class), |
||||
|
mock(TbServiceInfoProvider.class), |
||||
|
tenantProfileCache |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
// Regression test: publicUserSubscriptionsMap must be keyed by TenantId, not UserId(NULL_UUID).
|
||||
|
// With the old UserId(NULL_UUID) key, all tenants shared one global subscription counter.
|
||||
|
@Test |
||||
|
void processSubscription_publicUserSubscriptionsMap_isPerTenantNotGlobal() throws Exception { |
||||
|
int maxPublicSubscriptions = 2; |
||||
|
|
||||
|
TenantId tenant1 = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile1 = new TenantProfile(); |
||||
|
profile1.createDefaultTenantProfileData(); |
||||
|
profile1.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile1).given(tenantProfileCache).get(tenant1); |
||||
|
|
||||
|
TenantId tenant2 = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile2 = new TenantProfile(); |
||||
|
profile2.createDefaultTenantProfileData(); |
||||
|
profile2.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile2).given(tenantProfileCache).get(tenant2); |
||||
|
|
||||
|
// tenant1 fills up its quota
|
||||
|
for (int i = 0; i < maxPublicSubscriptions; i++) { |
||||
|
assertThat(service.processSubscription(mockPublicSessionRef(tenant1, "t1-session-" + i), subscriptionCmd(i))) |
||||
|
.as("tenant1 subscription %d should be accepted", i + 1) |
||||
|
.isTrue(); |
||||
|
} |
||||
|
|
||||
|
// tenant2 must have its own independent quota — this was the bug:
|
||||
|
// with UserId(NULL_UUID) as key all tenants shared one counter, so tenant2 would be blocked here
|
||||
|
for (int i = 0; i < maxPublicSubscriptions; i++) { |
||||
|
assertThat(service.processSubscription(mockPublicSessionRef(tenant2, "t2-session-" + i), subscriptionCmd(i))) |
||||
|
.as("tenant2 subscription %d should not be affected by tenant1's subscriptions", i + 1) |
||||
|
.isTrue(); |
||||
|
} |
||||
|
|
||||
|
// tenant1's (maxPublicSubscriptions + 1)-th subscription must be rejected
|
||||
|
assertThat(service.processSubscription(mockPublicSessionRef(tenant1, "t1-session-over"), subscriptionCmd(99))) |
||||
|
.as("tenant1 should be rejected after exceeding its limit") |
||||
|
.isFalse(); |
||||
|
|
||||
|
// Verify that publicUserSubscriptionsMap has separate entries per tenant
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
ConcurrentMap<TenantId, Set<String>> publicUserSubscriptionsMap = |
||||
|
(ConcurrentMap<TenantId, Set<String>>) ReflectionTestUtils.getField(service, "publicUserSubscriptionsMap"); |
||||
|
|
||||
|
assertThat(publicUserSubscriptionsMap).as("map should contain tenant1").containsKey(tenant1); |
||||
|
assertThat(publicUserSubscriptionsMap).as("map should contain tenant2").containsKey(tenant2); |
||||
|
assertThat(publicUserSubscriptionsMap).as("map must not have a single NULL_UUID entry for all tenants") |
||||
|
.doesNotContainKey(new TenantId(EntityId.NULL_UUID)); |
||||
|
|
||||
|
assertThat(publicUserSubscriptionsMap.get(tenant1)) |
||||
|
.as("tenant1 should have exactly %d subscriptions", maxPublicSubscriptions) |
||||
|
.hasSize(maxPublicSubscriptions); |
||||
|
assertThat(publicUserSubscriptionsMap.get(tenant2)) |
||||
|
.as("tenant2 should have exactly %d subscriptions", maxPublicSubscriptions) |
||||
|
.hasSize(maxPublicSubscriptions); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void processSubscription_publicUserSubscriptionsMap_subscriptionIdFormat() { |
||||
|
int maxPublicSubscriptions = 5; |
||||
|
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile = new TenantProfile(); |
||||
|
profile.createDefaultTenantProfileData(); |
||||
|
profile.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile).given(tenantProfileCache).get(tenantId); |
||||
|
|
||||
|
String sessionId = "my-session-id"; |
||||
|
int cmdId = 42; |
||||
|
WebSocketSessionRef sessionRef = mockPublicSessionRef(tenantId, sessionId); |
||||
|
service.processSubscription(sessionRef, subscriptionCmd(cmdId)); |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
ConcurrentMap<TenantId, Set<String>> publicUserSubscriptionsMap = |
||||
|
(ConcurrentMap<TenantId, Set<String>>) ReflectionTestUtils.getField(service, "publicUserSubscriptionsMap"); |
||||
|
|
||||
|
Set<String> subs = publicUserSubscriptionsMap.get(tenantId); |
||||
|
assertThat(subs).hasSize(1); |
||||
|
assertThat(subs.iterator().next()).isEqualTo("[" + sessionId + "]:[" + cmdId + "]"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void processSubscription_unsubscribe_removesEntryFromPublicUserSubscriptionsMap() { |
||||
|
int maxPublicSubscriptions = 5; |
||||
|
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile = new TenantProfile(); |
||||
|
profile.createDefaultTenantProfileData(); |
||||
|
profile.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile).given(tenantProfileCache).get(tenantId); |
||||
|
|
||||
|
String sessionId = "session-1"; |
||||
|
int cmdId = 1; |
||||
|
WebSocketSessionRef sessionRef = mockPublicSessionRef(tenantId, sessionId); |
||||
|
|
||||
|
service.processSubscription(sessionRef, subscriptionCmd(cmdId)); |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
ConcurrentMap<TenantId, Set<String>> publicUserSubscriptionsMap = |
||||
|
(ConcurrentMap<TenantId, Set<String>>) ReflectionTestUtils.getField(service, "publicUserSubscriptionsMap"); |
||||
|
assertThat(publicUserSubscriptionsMap.get(tenantId)).hasSize(1); |
||||
|
|
||||
|
AttributesSubscriptionCmd unsubCmd = subscriptionCmd(cmdId); |
||||
|
unsubCmd.setUnsubscribe(true); |
||||
|
service.processSubscription(sessionRef, unsubCmd); |
||||
|
|
||||
|
assertThat(publicUserSubscriptionsMap.get(tenantId)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void processSubscription_unsubscribe_freesSlotForNewSubscription() { |
||||
|
int maxPublicSubscriptions = 1; |
||||
|
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile = new TenantProfile(); |
||||
|
profile.createDefaultTenantProfileData(); |
||||
|
profile.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile).given(tenantProfileCache).get(tenantId); |
||||
|
|
||||
|
WebSocketSessionRef sessionRef = mockPublicSessionRef(tenantId, "session-1"); |
||||
|
service.processSubscription(sessionRef, subscriptionCmd(1)); |
||||
|
|
||||
|
// slot is full — second subscription on same session should be rejected
|
||||
|
assertThat(service.processSubscription(sessionRef, subscriptionCmd(2))).isFalse(); |
||||
|
|
||||
|
// unsubscribe cmd 1 to free the slot
|
||||
|
AttributesSubscriptionCmd unsubCmd = subscriptionCmd(1); |
||||
|
unsubCmd.setUnsubscribe(true); |
||||
|
service.processSubscription(sessionRef, unsubCmd); |
||||
|
|
||||
|
// now a new subscription should succeed
|
||||
|
assertThat(service.processSubscription(sessionRef, subscriptionCmd(3))) |
||||
|
.as("new subscription should succeed after unsubscribe freed the slot") |
||||
|
.isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void processSessionClose_removesAllSessionSubscriptionsFromPublicUserSubscriptionsMap() { |
||||
|
int maxPublicSubscriptions = 10; |
||||
|
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile = new TenantProfile(); |
||||
|
profile.createDefaultTenantProfileData(); |
||||
|
profile.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile).given(tenantProfileCache).get(tenantId); |
||||
|
|
||||
|
String sessionId = "closing-session"; |
||||
|
WebSocketSessionRef sessionRef = mockPublicSessionRef(tenantId, sessionId); |
||||
|
|
||||
|
service.processSubscription(sessionRef, subscriptionCmd(1)); |
||||
|
service.processSubscription(sessionRef, subscriptionCmd(2)); |
||||
|
service.processSubscription(sessionRef, subscriptionCmd(3)); |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
ConcurrentMap<TenantId, Set<String>> publicUserSubscriptionsMap = |
||||
|
(ConcurrentMap<TenantId, Set<String>>) ReflectionTestUtils.getField(service, "publicUserSubscriptionsMap"); |
||||
|
assertThat(publicUserSubscriptionsMap.get(tenantId)).hasSize(3); |
||||
|
|
||||
|
service.processSessionClose(sessionRef); |
||||
|
|
||||
|
assertThat(publicUserSubscriptionsMap.get(tenantId)).isEmpty(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void processSessionClose_onlyRemovesClosedSessionSubscriptions() { |
||||
|
int maxPublicSubscriptions = 10; |
||||
|
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); |
||||
|
TenantProfile profile = new TenantProfile(); |
||||
|
profile.createDefaultTenantProfileData(); |
||||
|
profile.getDefaultProfileConfiguration().setMaxWsSubscriptionsPerPublicUser(maxPublicSubscriptions); |
||||
|
willReturn(profile).given(tenantProfileCache).get(tenantId); |
||||
|
|
||||
|
WebSocketSessionRef session1 = mockPublicSessionRef(tenantId, "session-1"); |
||||
|
WebSocketSessionRef session2 = mockPublicSessionRef(tenantId, "session-2"); |
||||
|
|
||||
|
service.processSubscription(session1, subscriptionCmd(1)); |
||||
|
service.processSubscription(session1, subscriptionCmd(2)); |
||||
|
service.processSubscription(session2, subscriptionCmd(1)); |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
ConcurrentMap<TenantId, Set<String>> publicUserSubscriptionsMap = |
||||
|
(ConcurrentMap<TenantId, Set<String>>) ReflectionTestUtils.getField(service, "publicUserSubscriptionsMap"); |
||||
|
assertThat(publicUserSubscriptionsMap.get(tenantId)).hasSize(3); |
||||
|
|
||||
|
service.processSessionClose(session1); |
||||
|
|
||||
|
Set<String> remaining = publicUserSubscriptionsMap.get(tenantId); |
||||
|
assertThat(remaining).hasSize(1); |
||||
|
assertThat(remaining).allMatch(subId -> subId.startsWith("[session-2]")); |
||||
|
} |
||||
|
|
||||
|
private WebSocketSessionRef mockPublicSessionRef(TenantId tenantId, String sessionId) { |
||||
|
CustomerId customerId = new CustomerId(UUID.randomUUID()); |
||||
|
SecurityUser securityUser = mock(SecurityUser.class); |
||||
|
willReturn(tenantId).given(securityUser).getTenantId(); |
||||
|
willReturn(customerId).given(securityUser).getCustomerId(); |
||||
|
willReturn(new UserId(EntityId.NULL_UUID)).given(securityUser).getId(); |
||||
|
willReturn(true).given(securityUser).isCustomerUser(); |
||||
|
willReturn(new UserPrincipal(UserPrincipal.Type.PUBLIC_ID, customerId.toString())).given(securityUser).getUserPrincipal(); |
||||
|
|
||||
|
WebSocketSessionRef ref = mock(WebSocketSessionRef.class); |
||||
|
willReturn(securityUser).given(ref).getSecurityCtx(); |
||||
|
willReturn(sessionId).given(ref).getSessionId(); |
||||
|
return ref; |
||||
|
} |
||||
|
|
||||
|
private AttributesSubscriptionCmd subscriptionCmd(int cmdId) { |
||||
|
AttributesSubscriptionCmd cmd = new AttributesSubscriptionCmd(); |
||||
|
cmd.setCmdId(cmdId); |
||||
|
return cmd; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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,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,125 @@ |
|||||
|
/** |
||||
|
* 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. |
||||
|
*/ |
||||
|
syntax = "proto3"; |
||||
|
|
||||
|
option java_package = "org.thingsboard.server.gen.transport.coap"; |
||||
|
option java_outer_classname = "ConfigTypesProtos"; |
||||
|
|
||||
|
message ProtoCalibrationParametersRequest { |
||||
|
|
||||
|
/* Request details. * |
||||
|
* Bitmask: * |
||||
|
* - Bit 0:2 - Requested channel number, range: [1:6] * |
||||
|
* Status: Deprecated [06.10.00 - 06.xx.xx] */ |
||||
|
uint32 calibration_request = 1; |
||||
|
|
||||
|
/* Channel assignment - the sensor code. * |
||||
|
* Status: Deprecated [06.10.00 - 06.xx.xx] */ |
||||
|
uint32 channel_assignment = 2; |
||||
|
|
||||
|
/* Channel calibration parameters. * |
||||
|
* Up to 8 parameters supported. * |
||||
|
* If this field is empty, the sensor will send the current set of parameters. * |
||||
|
* Status: Deprecated [06.10.00 - 06.xx.xx] */ |
||||
|
repeated int32 parameters = 3; |
||||
|
} |
||||
|
|
||||
|
message ProtoOutputControlState { |
||||
|
|
||||
|
/* Channel index. * |
||||
|
* Range: [0:5] * |
||||
|
* Status: In use [06.13.00/06.21.00 - LATEST] */ |
||||
|
uint32 channel_index = 1; |
||||
|
|
||||
|
/* Channel output state: * |
||||
|
* - 1 - OFF * |
||||
|
* - 2 - ON * |
||||
|
* Status: In use [06.13.00/06.21.00 - LATEST] */ |
||||
|
uint32 channel_state = 2; |
||||
|
} |
||||
|
|
||||
|
enum BleAdvertisingPeriodMode { |
||||
|
|
||||
|
/* Invalid value. */ |
||||
|
BLE_ADVERTISING_PERIOD_MODE_UNSPECIFIED = 0; |
||||
|
|
||||
|
/* Default mode. * |
||||
|
* Bluetooth advertising interval is set to 1022.5ms or a lower value, based on the continuous measurement period. */ |
||||
|
BLE_ADVERTISING_PERIOD_MODE_DEFAULT = 1; |
||||
|
|
||||
|
/* Normal mode. * |
||||
|
* Uses the value configured by the user from the 'normal' field. */ |
||||
|
BLE_ADVERTISING_PERIOD_MODE_NORMAL = 2; |
||||
|
|
||||
|
/* Fast mode. * |
||||
|
* Uses the value configured by the user from the 'fast' field. */ |
||||
|
BLE_ADVERTISING_PERIOD_MODE_FAST = 3; |
||||
|
} |
||||
|
|
||||
|
message ProtoBleAdvertisingPeriod { |
||||
|
|
||||
|
/* Bluetooth advertising mode. * |
||||
|
* Status: In use [06.13.00/06.21.00 - LATEST] */ |
||||
|
BleAdvertisingPeriodMode mode = 1; |
||||
|
|
||||
|
/* Bluetooth advertising interval in normal mode, configured in steps of 0.625 ms. * |
||||
|
* Range: [32:16384] * |
||||
|
* Status: In use [06.13.00/06.21.00 - LATEST] */ |
||||
|
uint32 normal = 2; |
||||
|
|
||||
|
/* Bluetooth advertising interval in fast mode, configured in steps of 0.625 ms. * |
||||
|
* Range: [32:16384] * |
||||
|
* Status: In use [06.13.00/06.21.00 - LATEST] */ |
||||
|
uint32 fast = 3; |
||||
|
} |
||||
|
|
||||
|
enum AdvertisementManufacturerDataFormat { |
||||
|
|
||||
|
/* Invalid value. */ |
||||
|
ADVERTISEMENT_MANUFACTURER_DATA_FORMAT_UNSPECIFIED = 0; |
||||
|
|
||||
|
/* Advertisement manufacturer specific data format 3. */ |
||||
|
ADVERTISEMENT_MANUFACTURER_DATA_FORMAT_V3 = 1; |
||||
|
|
||||
|
/* Advertisement manufacturer specific data format 5. */ |
||||
|
ADVERTISEMENT_MANUFACTURER_DATA_FORMAT_V5 = 2; |
||||
|
} |
||||
|
|
||||
|
message ProtoNetworkSearch { |
||||
|
|
||||
|
/* Timing schema, if successful registration since the last reset. * |
||||
|
* Length: 6 items. * |
||||
|
* 1st - 6th item - Time in minutes. Range: [1:255]. * |
||||
|
* Status: In use [06.20.00 - LATEST] */ |
||||
|
repeated uint32 time_schema_last_registration_ok = 1; |
||||
|
|
||||
|
/* Timing schema, if no successful registration since the last reset. * |
||||
|
* Length: 6 items. * |
||||
|
* 1st - 6th item - Time in minutes. Range: [1:255]. * |
||||
|
* Status: In use [06.20.00 - LATEST] */ |
||||
|
repeated uint32 time_schema_last_registration_not_ok = 2; |
||||
|
|
||||
|
/* Disable base period in minutes. * |
||||
|
* Disable time = 'disable_period_base' * counter (from 1 to 'counter_max'). * |
||||
|
* Range: [1:255] * |
||||
|
* Status: In use [06.20.00 - LATEST] */ |
||||
|
uint32 disable_period_base = 3; |
||||
|
|
||||
|
/* Disable counter maximum. * |
||||
|
* Range: [1:255] * |
||||
|
* Status: In use [06.20.00 - LATEST] */ |
||||
|
uint32 counter_max = 4; |
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue