153 changed files with 2645 additions and 1165 deletions
@ -0,0 +1,56 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 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.controller; |
||||
|
|
||||
|
import org.junit.Test; |
||||
|
import org.stringtemplate.v4.ST; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.SaveDeviceWithCredentialsRequest; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentialsType; |
||||
|
|
||||
|
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
||||
|
|
||||
|
public abstract class BaseTelemetryControllerTest extends AbstractControllerTest { |
||||
|
|
||||
|
@Test |
||||
|
public void testConstraintValidator() throws Exception { |
||||
|
loginTenantAdmin(); |
||||
|
Device device = createDevice(); |
||||
|
String correctRequestBody = "{\"data\": \"value\"}"; |
||||
|
doPostAsync("/api/plugins/telemetry/" + device.getId() + "/SHARED_SCOPE", correctRequestBody, String.class, status().isOk()); |
||||
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/timeseries/smth", correctRequestBody, String.class, status().isOk()); |
||||
|
String invalidRequestBody = "{\"<object data=\\\"data:text/html,<script>alert(document)</script>\\\"></object>\": \"data\"}"; |
||||
|
doPostAsync("/api/plugins/telemetry/" + device.getId() + "/SHARED_SCOPE", invalidRequestBody, String.class, status().isBadRequest()); |
||||
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/timeseries/smth", invalidRequestBody, String.class, status().isBadRequest()); |
||||
|
} |
||||
|
|
||||
|
private Device createDevice() throws Exception { |
||||
|
String testToken = "TEST_TOKEN"; |
||||
|
|
||||
|
Device device = new Device(); |
||||
|
device.setName("My device"); |
||||
|
device.setType("default"); |
||||
|
|
||||
|
DeviceCredentials deviceCredentials = new DeviceCredentials(); |
||||
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); |
||||
|
deviceCredentials.setCredentialsId(testToken); |
||||
|
|
||||
|
SaveDeviceWithCredentialsRequest saveRequest = new SaveDeviceWithCredentialsRequest(device, deviceCredentials); |
||||
|
|
||||
|
return readResponse(doPost("/api/device-with-credentials", saveRequest).andExpect(status().isOk()), Device.class); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,266 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 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.device.provision; |
||||
|
|
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.assertj.core.api.Assertions; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.junit.runner.RunWith; |
||||
|
import org.springframework.boot.test.mock.mockito.MockBean; |
||||
|
import org.springframework.boot.test.mock.mockito.SpyBean; |
||||
|
import org.springframework.test.context.ContextConfiguration; |
||||
|
import org.springframework.test.context.junit4.SpringRunner; |
||||
|
import org.thingsboard.server.cluster.TbClusterService; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.DeviceProfile; |
||||
|
import org.thingsboard.server.common.data.DeviceProfileProvisionType; |
||||
|
import org.thingsboard.server.common.data.Tenant; |
||||
|
import org.thingsboard.server.common.data.device.credentials.ProvisionDeviceCredentialsData; |
||||
|
import org.thingsboard.server.common.data.device.profile.DeviceProfileData; |
||||
|
import org.thingsboard.server.common.data.device.profile.X509CertificateChainProvisionConfiguration; |
||||
|
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.TenantId; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentialsType; |
||||
|
import org.thingsboard.server.common.msg.EncryptionUtil; |
||||
|
import org.thingsboard.server.common.transport.util.SslUtil; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
import org.thingsboard.server.dao.audit.AuditLogService; |
||||
|
import org.thingsboard.server.dao.device.DeviceCredentialsService; |
||||
|
import org.thingsboard.server.dao.device.DeviceProfileService; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionFailedException; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionRequest; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionResponse; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionResponseStatus; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
import org.thingsboard.server.queue.TbQueueProducer; |
||||
|
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
||||
|
import org.thingsboard.server.queue.discovery.PartitionService; |
||||
|
import org.thingsboard.server.queue.provider.TbQueueProducerProvider; |
||||
|
import org.thingsboard.server.service.device.DeviceProvisionServiceImpl;;import java.io.IOException; |
||||
|
import java.nio.file.Files; |
||||
|
import java.nio.file.Paths; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
import java.util.regex.Matcher; |
||||
|
import java.util.regex.Pattern; |
||||
|
|
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@Slf4j |
||||
|
@RunWith(SpringRunner.class) |
||||
|
@ContextConfiguration(classes = DeviceProvisionServiceImpl.class) |
||||
|
public class DeviceProvisionServiceTest { |
||||
|
|
||||
|
@MockBean |
||||
|
protected TbQueueProducerProvider producerProvider; |
||||
|
@MockBean |
||||
|
protected TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg>> ruleEngineMsgProducer; |
||||
|
@MockBean |
||||
|
protected TbClusterService clusterService; |
||||
|
@MockBean |
||||
|
protected DeviceProfileService deviceProfileService; |
||||
|
@MockBean |
||||
|
protected DeviceService deviceService; |
||||
|
@MockBean |
||||
|
protected DeviceCredentialsService deviceCredentialsService; |
||||
|
@MockBean |
||||
|
protected AttributesService attributesService; |
||||
|
@MockBean |
||||
|
protected AuditLogService auditLogService; |
||||
|
@MockBean |
||||
|
protected PartitionService partitionService; |
||||
|
@SpyBean |
||||
|
DeviceProvisionServiceImpl service; |
||||
|
|
||||
|
private String[] chain; |
||||
|
|
||||
|
@Before |
||||
|
public void setUp() { |
||||
|
String filePath = "src/test/resources/provision/x509ChainProvisionTest.pem"; |
||||
|
try { |
||||
|
String certificateChain = Files.readString(Paths.get(filePath)); |
||||
|
certificateChain = certTrimNewLinesForChainInDeviceProfile(certificateChain); |
||||
|
chain = fetchLeafCertificateFromChain(certificateChain); |
||||
|
} catch (IOException e) { |
||||
|
throw new RuntimeException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
|
||||
|
@Test |
||||
|
public void provisionDeviceViaX509Certificate() { |
||||
|
var tenant = createTenant(); |
||||
|
var deviceProfile = createDeviceProfile(tenant.getId(), chain[1], true); |
||||
|
|
||||
|
var device = createDevice(tenant.getId(), deviceProfile.getId()); |
||||
|
when(deviceService.findDeviceByTenantIdAndName(any(), any())).thenReturn(device); |
||||
|
|
||||
|
var deviceCredentials = createDeviceCredentials(chain[0], device.getId()); |
||||
|
when(deviceCredentialsService.findDeviceCredentialsByDeviceId(any(), any())).thenReturn(deviceCredentials); |
||||
|
when(deviceCredentialsService.updateDeviceCredentials(any(), any())).thenReturn(deviceCredentials); |
||||
|
|
||||
|
ProvisionResponse response = service.provisionDeviceViaX509Chain(deviceProfile, createProvisionRequest(chain[0])); |
||||
|
|
||||
|
verify(deviceService, times(1)).findDeviceByTenantIdAndName(any(), any()); |
||||
|
verify(deviceCredentialsService, times(1)).findDeviceCredentialsByDeviceId(any(), any()); |
||||
|
verify(deviceCredentialsService, times(1)).updateDeviceCredentials(any(), any()); |
||||
|
|
||||
|
Assertions.assertThat(response.getResponseStatus()).isEqualTo(ProvisionResponseStatus.SUCCESS); |
||||
|
Assertions.assertThat(response.getDeviceCredentials()).isEqualTo(deviceCredentials); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void provisionDeviceWithIncorrectConfiguration() { |
||||
|
var tenant = createTenant(); |
||||
|
var deviceProfile = createDeviceProfile(tenant.getId(), chain[1], false); |
||||
|
|
||||
|
Assertions.assertThatThrownBy(() -> |
||||
|
service.provisionDeviceViaX509Chain(deviceProfile, createProvisionRequest(chain[0]))) |
||||
|
.isInstanceOf(ProvisionFailedException.class); |
||||
|
|
||||
|
verify(deviceService, times(1)).findDeviceByTenantIdAndName(any(), any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void matchDeviceNameFromX509CNCertificateByRegex() { |
||||
|
var tenant = createTenant(); |
||||
|
var deviceProfile = createDeviceProfile(tenant.getId(), chain[1], true); |
||||
|
X509CertificateChainProvisionConfiguration configuration = (X509CertificateChainProvisionConfiguration) deviceProfile.getProfileData().getProvisionConfiguration(); |
||||
|
String CN = getCNFromX509Certificate(chain[0]); |
||||
|
String deviceName = service.extractDeviceNameFromCNByRegEx(deviceProfile, CN, configuration.getCertificateRegExPattern()); |
||||
|
|
||||
|
Assertions.assertThat(deviceName).isNotBlank(); |
||||
|
Assertions.assertThat(deviceName).isEqualTo("deviceCertificate"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void matchDeviceNameFromCNByRegex() { |
||||
|
var CN = "DeviceA.company.com"; |
||||
|
var regex = "(.*)\\.company.com"; |
||||
|
var result = service.extractDeviceNameFromCNByRegEx(null, CN, regex); |
||||
|
Assertions.assertThat(result).isNotBlank(); |
||||
|
Assertions.assertThat(result).isEqualTo("DeviceA"); |
||||
|
|
||||
|
CN = "DeviceA@company.com"; |
||||
|
regex = "(.*)@company.com"; |
||||
|
result = service.extractDeviceNameFromCNByRegEx(null, CN, regex); |
||||
|
Assertions.assertThat(result).isNotBlank(); |
||||
|
Assertions.assertThat(result).isEqualTo("DeviceA"); |
||||
|
|
||||
|
CN = "prefixDeviceAsuffix@company.com"; |
||||
|
regex = "prefix(.*)suffix@company.com"; |
||||
|
result = service.extractDeviceNameFromCNByRegEx(null, CN, regex); |
||||
|
Assertions.assertThat(result).isNotBlank(); |
||||
|
Assertions.assertThat(result).isEqualTo("DeviceA"); |
||||
|
|
||||
|
CN = "prefixDeviceAsufix@company.com"; |
||||
|
regex = "prefix(.*)sufix@company.com"; |
||||
|
result = service.extractDeviceNameFromCNByRegEx(null, CN, regex); |
||||
|
Assertions.assertThat(result).isNotBlank(); |
||||
|
Assertions.assertThat(result).isEqualTo("DeviceA"); |
||||
|
|
||||
|
CN = "region.DeviceA.220423@company.com"; |
||||
|
regex = "\\D+\\.(.*)\\.\\d+@company.com"; |
||||
|
result = service.extractDeviceNameFromCNByRegEx(null, CN, regex); |
||||
|
Assertions.assertThat(result).isNotBlank(); |
||||
|
Assertions.assertThat(result).isEqualTo("DeviceA"); |
||||
|
} |
||||
|
|
||||
|
private DeviceProfile createDeviceProfile(TenantId tenantId, String certificateValue, boolean isAllowToCreateNewDevices) { |
||||
|
X509CertificateChainProvisionConfiguration provision = new X509CertificateChainProvisionConfiguration(); |
||||
|
provision.setProvisionDeviceSecret(certificateValue); |
||||
|
provision.setCertificateRegExPattern("([^@]+)"); |
||||
|
provision.setAllowCreateNewDevicesByX509Certificate(isAllowToCreateNewDevices); |
||||
|
|
||||
|
DeviceProfileData deviceProfileData = new DeviceProfileData(); |
||||
|
deviceProfileData.setProvisionConfiguration(provision); |
||||
|
|
||||
|
DeviceProfile deviceProfile = new DeviceProfile(); |
||||
|
deviceProfile.setId(new DeviceProfileId(UUID.randomUUID())); |
||||
|
deviceProfile.setProfileData(deviceProfileData); |
||||
|
deviceProfile.setProvisionDeviceKey(EncryptionUtil.getSha3Hash(certificateValue)); |
||||
|
deviceProfile.setProvisionType(DeviceProfileProvisionType.X509_CERTIFICATE_CHAIN); |
||||
|
deviceProfile.setTenantId(tenantId); |
||||
|
return deviceProfile; |
||||
|
} |
||||
|
|
||||
|
private Device createDevice(TenantId tenantId, DeviceProfileId deviceProfileId) { |
||||
|
Device device = new Device(); |
||||
|
device.setTenantId(tenantId); |
||||
|
device.setId(new DeviceId(UUID.randomUUID())); |
||||
|
device.setDeviceProfileId(deviceProfileId); |
||||
|
device.setCustomerId(new CustomerId(UUID.randomUUID())); |
||||
|
return device; |
||||
|
} |
||||
|
|
||||
|
private Tenant createTenant() { |
||||
|
Tenant tenant = new Tenant(); |
||||
|
tenant.setId(new TenantId(UUID.randomUUID())); |
||||
|
return tenant; |
||||
|
} |
||||
|
|
||||
|
private DeviceCredentials createDeviceCredentials(String certificateValue, DeviceId deviceId) { |
||||
|
DeviceCredentials deviceCredentials = new DeviceCredentials(); |
||||
|
deviceCredentials.setDeviceId(deviceId); |
||||
|
deviceCredentials.setCredentialsValue(certificateValue); |
||||
|
deviceCredentials.setCredentialsId(EncryptionUtil.getSha3Hash(certificateValue)); |
||||
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.X509_CERTIFICATE); |
||||
|
return deviceCredentials; |
||||
|
} |
||||
|
|
||||
|
private ProvisionRequest createProvisionRequest(String certificateValue) { |
||||
|
return new ProvisionRequest(null, DeviceCredentialsType.X509_CERTIFICATE, |
||||
|
new ProvisionDeviceCredentialsData(null, null, null, null, certificateValue), |
||||
|
null); |
||||
|
} |
||||
|
|
||||
|
public static String certTrimNewLinesForChainInDeviceProfile(String input) { |
||||
|
return input.replaceAll("\n", "") |
||||
|
.replaceAll("\r", "") |
||||
|
.replaceAll("-----BEGIN CERTIFICATE-----", "-----BEGIN CERTIFICATE-----\n") |
||||
|
.replaceAll("-----END CERTIFICATE-----", "\n-----END CERTIFICATE-----\n") |
||||
|
.trim(); |
||||
|
} |
||||
|
|
||||
|
private String[] fetchLeafCertificateFromChain(String value) { |
||||
|
List<String> chain = new ArrayList<>(); |
||||
|
String regex = "-----BEGIN CERTIFICATE-----\\s*.*?\\s*-----END CERTIFICATE-----"; |
||||
|
Pattern pattern = Pattern.compile(regex); |
||||
|
Matcher matcher = pattern.matcher(value); |
||||
|
while (matcher.find()) { |
||||
|
chain.add(matcher.group(0)); |
||||
|
} |
||||
|
return chain.toArray(new String[0]); |
||||
|
} |
||||
|
|
||||
|
private String getCNFromX509Certificate(String x509Value) { |
||||
|
try { |
||||
|
return SslUtil.parseCommonName(SslUtil.readCertFile(x509Value)); |
||||
|
} catch (Exception e) { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,211 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 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.transport; |
||||
|
|
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.junit.runner.RunWith; |
||||
|
import org.springframework.boot.test.mock.mockito.MockBean; |
||||
|
import org.springframework.boot.test.mock.mockito.SpyBean; |
||||
|
import org.springframework.test.context.ContextConfiguration; |
||||
|
import org.springframework.test.context.junit4.SpringRunner; |
||||
|
import org.thingsboard.server.cache.ota.OtaPackageDataCache; |
||||
|
import org.thingsboard.server.cluster.TbClusterService; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.DeviceProfile; |
||||
|
import org.thingsboard.server.common.data.DeviceProfileProvisionType; |
||||
|
import org.thingsboard.server.common.data.device.profile.DeviceProfileData; |
||||
|
import org.thingsboard.server.common.data.device.profile.X509CertificateChainProvisionConfiguration; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentialsType; |
||||
|
import org.thingsboard.server.common.msg.EncryptionUtil; |
||||
|
import org.thingsboard.server.dao.device.DeviceCredentialsService; |
||||
|
import org.thingsboard.server.dao.device.DeviceProfileService; |
||||
|
import org.thingsboard.server.dao.device.DeviceProvisionService; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionResponse; |
||||
|
import org.thingsboard.server.dao.device.provision.ProvisionResponseStatus; |
||||
|
import org.thingsboard.server.dao.ota.OtaPackageService; |
||||
|
import org.thingsboard.server.dao.queue.QueueService; |
||||
|
import org.thingsboard.server.dao.relation.RelationService; |
||||
|
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; |
||||
|
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
||||
|
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
||||
|
import org.thingsboard.server.service.executors.DbCallbackExecutorService; |
||||
|
import org.thingsboard.server.service.profile.TbDeviceProfileCache; |
||||
|
import org.thingsboard.server.service.resource.TbResourceService; |
||||
|
|
||||
|
import java.io.IOException; |
||||
|
import java.nio.file.Files; |
||||
|
import java.nio.file.Paths; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
import java.util.regex.Matcher; |
||||
|
import java.util.regex.Pattern; |
||||
|
|
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@Slf4j |
||||
|
@RunWith(SpringRunner.class) |
||||
|
@ContextConfiguration(classes = DefaultTransportApiService.class) |
||||
|
public class DefaultTransportApiServiceTest { |
||||
|
|
||||
|
@MockBean |
||||
|
protected TbDeviceProfileCache deviceProfileCache; |
||||
|
@MockBean |
||||
|
protected TbTenantProfileCache tenantProfileCache; |
||||
|
@MockBean |
||||
|
protected TbApiUsageStateService apiUsageStateService; |
||||
|
@MockBean |
||||
|
protected DeviceService deviceService; |
||||
|
@MockBean |
||||
|
protected DeviceProfileService deviceProfileService; |
||||
|
@MockBean |
||||
|
protected RelationService relationService; |
||||
|
@MockBean |
||||
|
protected DeviceCredentialsService deviceCredentialsService; |
||||
|
@MockBean |
||||
|
protected DbCallbackExecutorService dbCallbackExecutorService; |
||||
|
@MockBean |
||||
|
protected TbClusterService tbClusterService; |
||||
|
@MockBean |
||||
|
protected DataDecodingEncodingService dataDecodingEncodingService; |
||||
|
@MockBean |
||||
|
protected DeviceProvisionService deviceProvisionService; |
||||
|
@MockBean |
||||
|
protected TbResourceService resourceService; |
||||
|
@MockBean |
||||
|
protected OtaPackageService otaPackageService; |
||||
|
@MockBean |
||||
|
protected OtaPackageDataCache otaPackageDataCache; |
||||
|
@MockBean |
||||
|
protected QueueService queueService; |
||||
|
@SpyBean |
||||
|
DefaultTransportApiService service; |
||||
|
|
||||
|
private String certificateChain; |
||||
|
private String[] chain; |
||||
|
|
||||
|
@Before |
||||
|
public void setUp() { |
||||
|
|
||||
|
String filePath = "src/test/resources/provision/x509ChainProvisionTest.pem"; |
||||
|
try { |
||||
|
certificateChain = Files.readString(Paths.get(filePath)); |
||||
|
certificateChain = certTrimNewLinesForChainInDeviceProfile(certificateChain); |
||||
|
chain = fetchLeafCertificateFromChain(certificateChain); |
||||
|
} catch (IOException e) { |
||||
|
throw new RuntimeException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void validateExistingDeviceByX509CertificateStrategy() { |
||||
|
var device = createDevice(); |
||||
|
when(deviceService.findDeviceByIdAsync(any(), any())).thenReturn(Futures.immediateFuture(device)); |
||||
|
|
||||
|
var deviceCredentials = createDeviceCredentials(chain[0], device.getId()); |
||||
|
when(deviceCredentialsService.findDeviceCredentialsByCredentialsId(any())).thenReturn(deviceCredentials); |
||||
|
|
||||
|
service.validateOrCreateDeviceX509Certificate(certificateChain); |
||||
|
verify(deviceCredentialsService, times(1)).findDeviceCredentialsByCredentialsId(any()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void provisionDeviceX509Certificate() { |
||||
|
var deviceProfile = createDeviceProfile(chain[1]); |
||||
|
when(deviceProfileService.findDeviceProfileByProvisionDeviceKey(any())).thenReturn(deviceProfile); |
||||
|
|
||||
|
var device = createDevice(); |
||||
|
when(deviceService.findDeviceByTenantIdAndName(any(), any())).thenReturn(device); |
||||
|
when(deviceService.findDeviceByIdAsync(any(), any())).thenReturn(Futures.immediateFuture(device)); |
||||
|
|
||||
|
var deviceCredentials = createDeviceCredentials(chain[0], device.getId()); |
||||
|
when(deviceCredentialsService.findDeviceCredentialsByCredentialsId(any())).thenReturn(null); |
||||
|
when(deviceCredentialsService.updateDeviceCredentials(any(), any())).thenReturn(deviceCredentials); |
||||
|
|
||||
|
var provisionResponse = createProvisionResponse(deviceCredentials); |
||||
|
when(deviceProvisionService.provisionDeviceViaX509Chain(any(), any())).thenReturn(provisionResponse); |
||||
|
|
||||
|
service.validateOrCreateDeviceX509Certificate(certificateChain); |
||||
|
verify(deviceProfileService, times(1)).findDeviceProfileByProvisionDeviceKey(any()); |
||||
|
verify(deviceService, times(1)).findDeviceByIdAsync(any(), any()); |
||||
|
verify(deviceCredentialsService, times(1)).findDeviceCredentialsByCredentialsId(any()); |
||||
|
verify(deviceProvisionService, times(1)).provisionDeviceViaX509Chain(any(), any()); |
||||
|
} |
||||
|
|
||||
|
private DeviceProfile createDeviceProfile(String certificateValue) { |
||||
|
X509CertificateChainProvisionConfiguration provision = new X509CertificateChainProvisionConfiguration(); |
||||
|
provision.setProvisionDeviceSecret(certificateValue); |
||||
|
provision.setCertificateRegExPattern("([^@]+)"); |
||||
|
provision.setAllowCreateNewDevicesByX509Certificate(true); |
||||
|
|
||||
|
DeviceProfileData deviceProfileData = new DeviceProfileData(); |
||||
|
deviceProfileData.setProvisionConfiguration(provision); |
||||
|
|
||||
|
DeviceProfile deviceProfile = new DeviceProfile(); |
||||
|
deviceProfile.setProfileData(deviceProfileData); |
||||
|
deviceProfile.setProvisionDeviceKey(EncryptionUtil.getSha3Hash(certificateValue)); |
||||
|
deviceProfile.setProvisionType(DeviceProfileProvisionType.X509_CERTIFICATE_CHAIN); |
||||
|
return deviceProfile; |
||||
|
} |
||||
|
|
||||
|
private DeviceCredentials createDeviceCredentials(String certificateValue, DeviceId deviceId) { |
||||
|
DeviceCredentials deviceCredentials = new DeviceCredentials(); |
||||
|
deviceCredentials.setDeviceId(deviceId); |
||||
|
deviceCredentials.setCredentialsValue(certificateValue); |
||||
|
deviceCredentials.setCredentialsId(EncryptionUtil.getSha3Hash(certificateValue)); |
||||
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.X509_CERTIFICATE); |
||||
|
return deviceCredentials; |
||||
|
} |
||||
|
|
||||
|
private Device createDevice() { |
||||
|
Device device = new Device(); |
||||
|
device.setId(new DeviceId(UUID.randomUUID())); |
||||
|
return device; |
||||
|
} |
||||
|
|
||||
|
private ProvisionResponse createProvisionResponse(DeviceCredentials deviceCredentials) { |
||||
|
return new ProvisionResponse(deviceCredentials, ProvisionResponseStatus.SUCCESS); |
||||
|
} |
||||
|
|
||||
|
public static String certTrimNewLinesForChainInDeviceProfile(String input) { |
||||
|
return input.replaceAll("\n", "") |
||||
|
.replaceAll("\r", "") |
||||
|
.replaceAll("-----BEGIN CERTIFICATE-----", "-----BEGIN CERTIFICATE-----\n") |
||||
|
.replaceAll("-----END CERTIFICATE-----", "\n-----END CERTIFICATE-----\n") |
||||
|
.trim(); |
||||
|
} |
||||
|
|
||||
|
private String[] fetchLeafCertificateFromChain(String value) { |
||||
|
List<String> chain = new ArrayList<>(); |
||||
|
String regex = "-----BEGIN CERTIFICATE-----\\s*.*?\\s*-----END CERTIFICATE-----"; |
||||
|
Pattern pattern = Pattern.compile(regex); |
||||
|
Matcher matcher = pattern.matcher(value); |
||||
|
while (matcher.find()) { |
||||
|
chain.add(matcher.group(0)); |
||||
|
} |
||||
|
return chain.toArray(new String[0]); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,36 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 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.system; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.Test; |
||||
|
import org.junit.jupiter.api.Assertions; |
||||
|
import org.springframework.util.ClassUtils; |
||||
|
import org.springframework.web.client.RestTemplate; |
||||
|
|
||||
|
|
||||
|
@Slf4j |
||||
|
public class RestTemplateConvertersTest { |
||||
|
|
||||
|
@Test |
||||
|
public void testJacksonXmlConverter() { |
||||
|
ClassLoader classLoader = RestTemplate.class.getClassLoader(); |
||||
|
boolean jackson2XmlPresent = ClassUtils.isPresent("com.fasterxml.jackson.dataformat.xml.XmlMapper", classLoader); |
||||
|
Assertions.assertFalse(jackson2XmlPresent, "XmlMapper must not be present in classpath, please, exclude \"jackson-dataformat-xml\" dependency!"); |
||||
|
//If this xml mapper will be present in classpath then we will get "Unsupported Media Type" in RestTemplate
|
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
-----BEGIN CERTIFICATE----- |
||||
|
MIICMTCCAdegAwIBAgIUI9dBuwN6pTtK6uZ03rkiCwV4wEYwCgYIKoZIzj0EAwIw |
||||
|
bjELMAkGA1UEBhMCVVMxETAPBgNVBAgMCE5ldyBZb3JrMRowGAYDVQQKDBFUaGlu |
||||
|
Z3NCb2FyZCwgSW5jLjEwMC4GA1UEAwwnZGV2aWNlQ2VydGlmaWNhdGVAWDUwOVBy |
||||
|
b3Zpc2lvblN0cmF0ZWd5MB4XDTIzMDMyOTE0NTYxN1oXDTI0MDMyODE0NTYxN1ow |
||||
|
bjELMAkGA1UEBhMCVVMxETAPBgNVBAgMCE5ldyBZb3JrMRowGAYDVQQKDBFUaGlu |
||||
|
Z3NCb2FyZCwgSW5jLjEwMC4GA1UEAwwnZGV2aWNlQ2VydGlmaWNhdGVAWDUwOVBy |
||||
|
b3Zpc2lvblN0cmF0ZWd5MFkwEwYHKoZIzj0CAQYIKoZIzj0DAQcDQgAE9Zo791qK |
||||
|
QiGNBm11r4ZGxh+w+ossZL3xc46ufq5QckQHP7zkD2XDAcmP5GvdkM1sBFN9AWaC |
||||
|
kQfNnWmfERsOOKNTMFEwHQYDVR0OBBYEFFFc5uyCyglQoZiKhzXzMcQ3BKORMB8G |
||||
|
A1UdIwQYMBaAFFFc5uyCyglQoZiKhzXzMcQ3BKORMA8GA1UdEwEB/wQFMAMBAf8w |
||||
|
CgYIKoZIzj0EAwIDSAAwRQIhANbA9CuhoOifZMMmqkpuld+65CR+ItKdXeRAhLMZ |
||||
|
uccuAiB0FSQB34zMutXrZj1g8Gl5OkE7YryFHbei1z0SveHR8g== |
||||
|
-----END CERTIFICATE----- |
||||
|
-----BEGIN CERTIFICATE----- |
||||
|
MIICMTCCAdegAwIBAgIUUEKxS9hTz4l+oLUMF0LV6TC/gCIwCgYIKoZIzj0EAwIw |
||||
|
bjELMAkGA1UEBhMCVVMxETAPBgNVBAgMCE5ldyBZb3JrMRowGAYDVQQKDBFUaGlu |
||||
|
Z3NCb2FyZCwgSW5jLjEwMC4GA1UEAwwnZGV2aWNlUHJvZmlsZUNlcnRAWDUwOVBy |
||||
|
b3Zpc2lvblN0cmF0ZWd5MB4XDTIzMDMyOTE0NTczNloXDTI0MDMyODE0NTczNlow |
||||
|
bjELMAkGA1UEBhMCVVMxETAPBgNVBAgMCE5ldyBZb3JrMRowGAYDVQQKDBFUaGlu |
||||
|
Z3NCb2FyZCwgSW5jLjEwMC4GA1UEAwwnZGV2aWNlUHJvZmlsZUNlcnRAWDUwOVBy |
||||
|
b3Zpc2lvblN0cmF0ZWd5MFkwEwYHKoZIzj0CAQYIKoZIzj0DAQcDQgAECMlWO72k |
||||
|
rDoUL9FQjUmSCetkhaEGJUfQkdSfkLSNa0GyAEIMbfmzI4zITeapunu4rGet3EMy |
||||
|
LydQzuQanBicp6NTMFEwHQYDVR0OBBYEFHpZ78tPnztNii4Da/yCw6mhEIL3MB8G |
||||
|
A1UdIwQYMBaAFHpZ78tPnztNii4Da/yCw6mhEIL3MA8GA1UdEwEB/wQFMAMBAf8w |
||||
|
CgYIKoZIzj0EAwIDSAAwRQIgJ7qyMFqNcwSYkH6o+UlQXzLWfwZbNjVk+aR7foAZ |
||||
|
NGsCIQDsd7v3WQIGHiArfZeDs1DLEDuV/2h6L+ZNoGNhEKL+1A== |
||||
|
-----END CERTIFICATE----- |
||||
@ -0,0 +1,36 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 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.device.profile; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
import org.thingsboard.server.common.data.DeviceProfileProvisionType; |
||||
|
|
||||
|
@Data |
||||
|
@NoArgsConstructor |
||||
|
public class X509CertificateChainProvisionConfiguration implements DeviceProfileProvisionConfiguration { |
||||
|
|
||||
|
private String provisionDeviceSecret; |
||||
|
private String certificateRegExPattern; |
||||
|
private boolean allowCreateNewDevicesByX509Certificate; |
||||
|
|
||||
|
@Override |
||||
|
public DeviceProfileProvisionType getType() { |
||||
|
return DeviceProfileProvisionType.X509_CERTIFICATE_CHAIN; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,35 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.dao.util; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.kv.KvEntry; |
||||
|
import org.thingsboard.server.dao.exception.IncorrectParameterException; |
||||
|
import org.thingsboard.server.dao.service.ConstraintValidator; |
||||
|
|
||||
|
import java.util.List; |
||||
|
|
||||
|
public class KvUtils { |
||||
|
public static void validate(List<? extends KvEntry> tsKvEntries) { |
||||
|
tsKvEntries.forEach(KvUtils::validate); |
||||
|
} |
||||
|
|
||||
|
public static void validate(KvEntry tsKvEntry) { |
||||
|
if (tsKvEntry == null) { |
||||
|
throw new IncorrectParameterException("Key value entry can't be null"); |
||||
|
} |
||||
|
ConstraintValidator.validateFields(tsKvEntry); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,50 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2023 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.dao.service; |
||||
|
|
||||
|
import org.junit.Assert; |
||||
|
import org.junit.jupiter.api.Assertions; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
import org.thingsboard.server.dao.exception.DataValidationException; |
||||
|
|
||||
|
class ConstraintValidatorTest { |
||||
|
|
||||
|
private static final int MIN_IN_MS = 60000; |
||||
|
private static final int _1M = 1_000_000; |
||||
|
|
||||
|
@Test |
||||
|
void validateFields() { |
||||
|
StringDataEntry stringDataEntryValid = new StringDataEntry("key", "value"); |
||||
|
StringDataEntry stringDataEntryInvalid1 = new StringDataEntry("<object type=\"text/html\"><script>alert(document)</script></object>", "value"); |
||||
|
|
||||
|
Assert.assertThrows(DataValidationException.class, () -> ConstraintValidator.validateFields(stringDataEntryInvalid1)); |
||||
|
ConstraintValidator.validateFields(stringDataEntryValid); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void validatePerMinute() { |
||||
|
StringDataEntry stringDataEntryValid = new StringDataEntry("key", "value"); |
||||
|
|
||||
|
long start = System.currentTimeMillis(); |
||||
|
for (int i = 0; i < _1M; i++) { |
||||
|
ConstraintValidator.validateFields(stringDataEntryValid); |
||||
|
} |
||||
|
long end = System.currentTimeMillis(); |
||||
|
|
||||
|
Assertions.assertTrue(MIN_IN_MS > end - start); |
||||
|
} |
||||
|
} |
||||
@ -1,3 +0,0 @@ |
|||||
tb.baseUrl=http://localhost:8080 |
|
||||
tb.baseUiUrl=http://localhost:8080 |
|
||||
tb.wsUrl=ws://localhost:8080 |
|
||||
@ -1,402 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2023 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.rule.engine.action; |
|
||||
|
|
||||
import com.fasterxml.jackson.databind.node.ArrayNode; |
|
||||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.junit.jupiter.api.AfterEach; |
|
||||
import org.junit.jupiter.api.Assertions; |
|
||||
import org.junit.jupiter.api.BeforeEach; |
|
||||
import org.junit.jupiter.api.Test; |
|
||||
import org.mockito.ArgumentCaptor; |
|
||||
import org.mockito.ArgumentMatchers; |
|
||||
import org.mockito.stubbing.Answer; |
|
||||
import org.thingsboard.common.util.JacksonUtil; |
|
||||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
||||
import org.thingsboard.rule.engine.api.TbContext; |
|
||||
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|
||||
import org.thingsboard.rule.engine.api.TbNodeException; |
|
||||
import org.thingsboard.rule.engine.api.TbRelationTypes; |
|
||||
import org.thingsboard.rule.engine.deduplication.DeduplicationStrategy; |
|
||||
import org.thingsboard.rule.engine.deduplication.TbMsgDeduplicationNode; |
|
||||
import org.thingsboard.rule.engine.deduplication.TbMsgDeduplicationNodeConfiguration; |
|
||||
import org.thingsboard.server.common.data.id.DeviceId; |
|
||||
import org.thingsboard.server.common.data.id.EntityId; |
|
||||
import org.thingsboard.server.common.data.id.RuleNodeId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.TbMsg; |
|
||||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
||||
import org.thingsboard.server.common.msg.session.SessionMsgType; |
|
||||
|
|
||||
import java.util.ArrayList; |
|
||||
import java.util.List; |
|
||||
import java.util.Random; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.CountDownLatch; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.ScheduledExecutorService; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
import java.util.concurrent.atomic.AtomicLong; |
|
||||
import java.util.function.Consumer; |
|
||||
|
|
||||
import static org.mockito.ArgumentMatchers.any; |
|
||||
import static org.mockito.ArgumentMatchers.eq; |
|
||||
import static org.mockito.ArgumentMatchers.isNull; |
|
||||
import static org.mockito.ArgumentMatchers.nullable; |
|
||||
import static org.mockito.Mockito.doAnswer; |
|
||||
import static org.mockito.Mockito.mock; |
|
||||
import static org.mockito.Mockito.spy; |
|
||||
import static org.mockito.Mockito.times; |
|
||||
import static org.mockito.Mockito.verify; |
|
||||
import static org.mockito.Mockito.when; |
|
||||
|
|
||||
@Slf4j |
|
||||
public class TbMsgDeduplicationNodeTest { |
|
||||
|
|
||||
private static final String MAIN_QUEUE_NAME = "Main"; |
|
||||
private static final String HIGH_PRIORITY_QUEUE_NAME = "HighPriority"; |
|
||||
private static final String TB_MSG_DEDUPLICATION_TIMEOUT_MSG = "TbMsgDeduplicationNodeMsg"; |
|
||||
|
|
||||
private TbContext ctx; |
|
||||
|
|
||||
private final ThingsBoardThreadFactory factory = ThingsBoardThreadFactory.forName("de-duplication-node-test"); |
|
||||
private final ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor(factory); |
|
||||
private final int deduplicationInterval = 1; |
|
||||
|
|
||||
private TenantId tenantId; |
|
||||
|
|
||||
private TbMsgDeduplicationNode node; |
|
||||
private TbMsgDeduplicationNodeConfiguration config; |
|
||||
private TbNodeConfiguration nodeConfiguration; |
|
||||
|
|
||||
private CountDownLatch awaitTellSelfLatch; |
|
||||
|
|
||||
@BeforeEach |
|
||||
public void init() throws TbNodeException { |
|
||||
ctx = mock(TbContext.class); |
|
||||
|
|
||||
tenantId = TenantId.fromUUID(UUID.randomUUID()); |
|
||||
RuleNodeId ruleNodeId = new RuleNodeId(UUID.randomUUID()); |
|
||||
|
|
||||
when(ctx.getSelfId()).thenReturn(ruleNodeId); |
|
||||
when(ctx.getTenantId()).thenReturn(tenantId); |
|
||||
|
|
||||
doAnswer((Answer<TbMsg>) invocationOnMock -> { |
|
||||
String type = (String) (invocationOnMock.getArguments())[1]; |
|
||||
EntityId originator = (EntityId) (invocationOnMock.getArguments())[2]; |
|
||||
TbMsgMetaData metaData = (TbMsgMetaData) (invocationOnMock.getArguments())[3]; |
|
||||
String data = (String) (invocationOnMock.getArguments())[4]; |
|
||||
return TbMsg.newMsg(type, originator, metaData.copy(), data); |
|
||||
}).when(ctx).newMsg(isNull(), eq(TB_MSG_DEDUPLICATION_TIMEOUT_MSG), nullable(EntityId.class), any(TbMsgMetaData.class), any(String.class)); |
|
||||
node = spy(new TbMsgDeduplicationNode()); |
|
||||
config = new TbMsgDeduplicationNodeConfiguration().defaultConfiguration(); |
|
||||
} |
|
||||
|
|
||||
private void invokeTellSelf(int maxNumberOfInvocation) { |
|
||||
invokeTellSelf(maxNumberOfInvocation, false, 0); |
|
||||
} |
|
||||
|
|
||||
private void invokeTellSelf(int maxNumberOfInvocation, boolean delayScheduleTimeout, int delayMultiplier) { |
|
||||
AtomicLong scheduleTimeout = new AtomicLong(deduplicationInterval); |
|
||||
AtomicInteger scheduleCount = new AtomicInteger(0); |
|
||||
doAnswer((Answer<Void>) invocationOnMock -> { |
|
||||
scheduleCount.getAndIncrement(); |
|
||||
if (scheduleCount.get() <= maxNumberOfInvocation) { |
|
||||
TbMsg msg = (TbMsg) (invocationOnMock.getArguments())[0]; |
|
||||
executorService.schedule(() -> { |
|
||||
try { |
|
||||
node.onMsg(ctx, msg); |
|
||||
awaitTellSelfLatch.countDown(); |
|
||||
} catch (ExecutionException | InterruptedException | TbNodeException e) { |
|
||||
log.error("Failed to execute tellSelf method call due to: ", e); |
|
||||
} |
|
||||
}, scheduleTimeout.get(), TimeUnit.SECONDS); |
|
||||
if (delayScheduleTimeout) { |
|
||||
scheduleTimeout.set(scheduleTimeout.get() * delayMultiplier); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
return null; |
|
||||
}).when(ctx).tellSelf(ArgumentMatchers.any(TbMsg.class), ArgumentMatchers.anyLong()); |
|
||||
} |
|
||||
|
|
||||
@AfterEach |
|
||||
public void destroy() { |
|
||||
executorService.shutdown(); |
|
||||
node.destroy(); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void given_100_messages_strategy_first_then_verifyOutput() throws TbNodeException, ExecutionException, InterruptedException { |
|
||||
int wantedNumberOfTellSelfInvocation = 2; |
|
||||
int msgCount = 100; |
|
||||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|
||||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|
||||
|
|
||||
config.setInterval(deduplicationInterval); |
|
||||
config.setMaxPendingMsgs(msgCount); |
|
||||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
||||
node.init(ctx, nodeConfiguration); |
|
||||
|
|
||||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
||||
long currentTimeMillis = System.currentTimeMillis(); |
|
||||
|
|
||||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|
||||
for (TbMsg msg : inputMsgs) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
|
|
||||
TbMsg msgToReject = createMsg(deviceId, inputMsgs.get(inputMsgs.size() - 1).getMetaDataTs() + 2); |
|
||||
node.onMsg(ctx, msgToReject); |
|
||||
|
|
||||
awaitTellSelfLatch.await(); |
|
||||
|
|
||||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
||||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|
||||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|
||||
|
|
||||
verify(ctx, times(msgCount)).ack(any()); |
|
||||
verify(ctx, times(1)).tellFailure(eq(msgToReject), any()); |
|
||||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any()); |
|
||||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|
||||
Assertions.assertEquals(inputMsgs.get(0), newMsgCaptor.getValue()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void given_100_messages_strategy_last_then_verifyOutput() throws TbNodeException, ExecutionException, InterruptedException { |
|
||||
int wantedNumberOfTellSelfInvocation = 2; |
|
||||
int msgCount = 100; |
|
||||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|
||||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|
||||
|
|
||||
config.setStrategy(DeduplicationStrategy.LAST); |
|
||||
config.setInterval(deduplicationInterval); |
|
||||
config.setMaxPendingMsgs(msgCount); |
|
||||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
||||
node.init(ctx, nodeConfiguration); |
|
||||
|
|
||||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
||||
long currentTimeMillis = System.currentTimeMillis(); |
|
||||
|
|
||||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|
||||
TbMsg msgWithLatestTs = getMsgWithLatestTs(inputMsgs); |
|
||||
|
|
||||
for (TbMsg msg : inputMsgs) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
|
|
||||
TbMsg msgToReject = createMsg(deviceId, inputMsgs.get(inputMsgs.size() - 1).getMetaDataTs() + 2); |
|
||||
node.onMsg(ctx, msgToReject); |
|
||||
|
|
||||
awaitTellSelfLatch.await(); |
|
||||
|
|
||||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
||||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|
||||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|
||||
|
|
||||
verify(ctx, times(msgCount)).ack(any()); |
|
||||
verify(ctx, times(1)).tellFailure(eq(msgToReject), any()); |
|
||||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any()); |
|
||||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|
||||
Assertions.assertEquals(msgWithLatestTs, newMsgCaptor.getValue()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void given_100_messages_strategy_all_then_verifyOutput() throws TbNodeException, ExecutionException, InterruptedException { |
|
||||
int wantedNumberOfTellSelfInvocation = 2; |
|
||||
int msgCount = 100; |
|
||||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|
||||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|
||||
|
|
||||
config.setInterval(deduplicationInterval); |
|
||||
config.setStrategy(DeduplicationStrategy.ALL); |
|
||||
config.setOutMsgType(SessionMsgType.POST_ATTRIBUTES_REQUEST.name()); |
|
||||
config.setQueueName(HIGH_PRIORITY_QUEUE_NAME); |
|
||||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
||||
node.init(ctx, nodeConfiguration); |
|
||||
|
|
||||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
||||
long currentTimeMillis = System.currentTimeMillis(); |
|
||||
|
|
||||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|
||||
for (TbMsg msg : inputMsgs) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
|
|
||||
awaitTellSelfLatch.await(); |
|
||||
|
|
||||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
||||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|
||||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|
||||
|
|
||||
verify(ctx, times(msgCount)).ack(any()); |
|
||||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any()); |
|
||||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|
||||
|
|
||||
Assertions.assertEquals(1, newMsgCaptor.getAllValues().size()); |
|
||||
TbMsg outMessage = newMsgCaptor.getAllValues().get(0); |
|
||||
Assertions.assertEquals(getMergedData(inputMsgs), outMessage.getData()); |
|
||||
Assertions.assertEquals(deviceId, outMessage.getOriginator()); |
|
||||
Assertions.assertEquals(config.getOutMsgType(), outMessage.getType()); |
|
||||
Assertions.assertEquals(config.getQueueName(), outMessage.getQueueName()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void given_100_messages_strategy_all_then_verifyOutput_2_packs() throws TbNodeException, ExecutionException, InterruptedException { |
|
||||
int wantedNumberOfTellSelfInvocation = 2; |
|
||||
int msgCount = 100; |
|
||||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|
||||
invokeTellSelf(wantedNumberOfTellSelfInvocation, true, 3); |
|
||||
|
|
||||
config.setInterval(deduplicationInterval); |
|
||||
config.setStrategy(DeduplicationStrategy.ALL); |
|
||||
config.setOutMsgType(SessionMsgType.POST_ATTRIBUTES_REQUEST.name()); |
|
||||
config.setQueueName(HIGH_PRIORITY_QUEUE_NAME); |
|
||||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
||||
node.init(ctx, nodeConfiguration); |
|
||||
|
|
||||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
||||
long currentTimeMillis = System.currentTimeMillis(); |
|
||||
|
|
||||
List<TbMsg> firstMsgPack = getTbMsgs(deviceId, msgCount / 2, currentTimeMillis, 500); |
|
||||
for (TbMsg msg : firstMsgPack) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
long firstPackDeduplicationPackEndTs = firstMsgPack.get(0).getMetaDataTs() + TimeUnit.SECONDS.toMillis(deduplicationInterval); |
|
||||
|
|
||||
List<TbMsg> secondMsgPack = getTbMsgs(deviceId, msgCount / 2, firstPackDeduplicationPackEndTs, 500); |
|
||||
for (TbMsg msg : secondMsgPack) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
|
|
||||
awaitTellSelfLatch.await(); |
|
||||
|
|
||||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
||||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|
||||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|
||||
|
|
||||
verify(ctx, times(msgCount)).ack(any()); |
|
||||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any()); |
|
||||
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|
||||
|
|
||||
List<TbMsg> resultMsgs = newMsgCaptor.getAllValues(); |
|
||||
Assertions.assertEquals(2, resultMsgs.size()); |
|
||||
|
|
||||
TbMsg firstMsg = resultMsgs.get(0); |
|
||||
Assertions.assertEquals(getMergedData(firstMsgPack), firstMsg.getData()); |
|
||||
Assertions.assertEquals(deviceId, firstMsg.getOriginator()); |
|
||||
Assertions.assertEquals(config.getOutMsgType(), firstMsg.getType()); |
|
||||
Assertions.assertEquals(config.getQueueName(), firstMsg.getQueueName()); |
|
||||
|
|
||||
TbMsg secondMsg = resultMsgs.get(1); |
|
||||
Assertions.assertEquals(getMergedData(secondMsgPack), secondMsg.getData()); |
|
||||
Assertions.assertEquals(deviceId, secondMsg.getOriginator()); |
|
||||
Assertions.assertEquals(config.getOutMsgType(), secondMsg.getType()); |
|
||||
Assertions.assertEquals(config.getQueueName(), secondMsg.getQueueName()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void given_100_messages_strategy_last_then_verifyOutput_2_packs() throws TbNodeException, ExecutionException, InterruptedException { |
|
||||
int wantedNumberOfTellSelfInvocation = 2; |
|
||||
int msgCount = 100; |
|
||||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|
||||
invokeTellSelf(wantedNumberOfTellSelfInvocation, true, 3); |
|
||||
|
|
||||
config.setInterval(deduplicationInterval); |
|
||||
config.setStrategy(DeduplicationStrategy.LAST); |
|
||||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
||||
node.init(ctx, nodeConfiguration); |
|
||||
|
|
||||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
||||
long currentTimeMillis = System.currentTimeMillis(); |
|
||||
|
|
||||
List<TbMsg> firstMsgPack = getTbMsgs(deviceId, msgCount / 2, currentTimeMillis, 500); |
|
||||
for (TbMsg msg : firstMsgPack) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
long firstPackDeduplicationPackEndTs = firstMsgPack.get(0).getMetaDataTs() + TimeUnit.SECONDS.toMillis(deduplicationInterval); |
|
||||
TbMsg msgWithLatestTsInFirstPack = getMsgWithLatestTs(firstMsgPack); |
|
||||
|
|
||||
List<TbMsg> secondMsgPack = getTbMsgs(deviceId, msgCount / 2, firstPackDeduplicationPackEndTs, 500); |
|
||||
for (TbMsg msg : secondMsgPack) { |
|
||||
node.onMsg(ctx, msg); |
|
||||
} |
|
||||
TbMsg msgWithLatestTsInSecondPack = getMsgWithLatestTs(secondMsgPack); |
|
||||
|
|
||||
awaitTellSelfLatch.await(); |
|
||||
|
|
||||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
||||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|
||||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|
||||
|
|
||||
verify(ctx, times(msgCount)).ack(any()); |
|
||||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any()); |
|
||||
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|
||||
|
|
||||
List<TbMsg> resultMsgs = newMsgCaptor.getAllValues(); |
|
||||
Assertions.assertEquals(2, resultMsgs.size()); |
|
||||
Assertions.assertTrue(resultMsgs.contains(msgWithLatestTsInFirstPack)); |
|
||||
Assertions.assertTrue(resultMsgs.contains(msgWithLatestTsInSecondPack)); |
|
||||
} |
|
||||
|
|
||||
private TbMsg getMsgWithLatestTs(List<TbMsg> firstMsgPack) { |
|
||||
int indexOfLastMsgInArray = firstMsgPack.size() - 1; |
|
||||
int indexToSetMaxTs = new Random().nextInt(indexOfLastMsgInArray) + 1; |
|
||||
TbMsg currentMaxTsMsg = firstMsgPack.get(indexOfLastMsgInArray); |
|
||||
TbMsg newLastMsgOfArray = firstMsgPack.get(indexToSetMaxTs); |
|
||||
firstMsgPack.set(indexOfLastMsgInArray, newLastMsgOfArray); |
|
||||
firstMsgPack.set(indexToSetMaxTs, currentMaxTsMsg); |
|
||||
return currentMaxTsMsg; |
|
||||
} |
|
||||
|
|
||||
private List<TbMsg> getTbMsgs(DeviceId deviceId, int msgCount, long currentTimeMillis, int initTsStep) { |
|
||||
List<TbMsg> inputMsgs = new ArrayList<>(); |
|
||||
var ts = currentTimeMillis + initTsStep; |
|
||||
for (int i = 0; i < msgCount; i++) { |
|
||||
inputMsgs.add(createMsg(deviceId, ts)); |
|
||||
ts += 2; |
|
||||
} |
|
||||
return inputMsgs; |
|
||||
} |
|
||||
|
|
||||
private TbMsg createMsg(DeviceId deviceId, long ts) { |
|
||||
ObjectNode dataNode = JacksonUtil.newObjectNode(); |
|
||||
dataNode.put("deviceId", deviceId.getId().toString()); |
|
||||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|
||||
metaData.putValue("ts", String.valueOf(ts)); |
|
||||
return TbMsg.newMsg( |
|
||||
MAIN_QUEUE_NAME, |
|
||||
SessionMsgType.POST_TELEMETRY_REQUEST.name(), |
|
||||
deviceId, |
|
||||
metaData, |
|
||||
JacksonUtil.toString(dataNode)); |
|
||||
} |
|
||||
|
|
||||
private String getMergedData(List<TbMsg> msgs) { |
|
||||
ArrayNode mergedData = JacksonUtil.OBJECT_MAPPER.createArrayNode(); |
|
||||
msgs.forEach(msg -> { |
|
||||
ObjectNode msgNode = JacksonUtil.newObjectNode(); |
|
||||
msgNode.set("msg", JacksonUtil.toJsonNode(msg.getData())); |
|
||||
msgNode.set("metadata", JacksonUtil.valueToTree(msg.getMetaData().getData())); |
|
||||
mergedData.add(msgNode); |
|
||||
}); |
|
||||
return JacksonUtil.toString(mergedData); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue