136 changed files with 7448 additions and 3881 deletions
@ -0,0 +1,247 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.queue; |
||||
|
|
||||
|
import com.google.common.collect.Sets; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
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.id.QueueId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.queue.Queue; |
||||
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
||||
|
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.NotificationsTopicService; |
||||
|
import org.thingsboard.server.queue.discovery.PartitionService; |
||||
|
import org.thingsboard.server.queue.provider.TbQueueProducerProvider; |
||||
|
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
||||
|
import org.thingsboard.server.service.gateway_device.GatewayNotificationsService; |
||||
|
import org.thingsboard.server.service.profile.TbAssetProfileCache; |
||||
|
import org.thingsboard.server.service.profile.TbDeviceProfileCache; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.ArgumentMatchers.isNull; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@Slf4j |
||||
|
@RunWith(SpringRunner.class) |
||||
|
@ContextConfiguration(classes = DefaultTbClusterService.class) |
||||
|
public class DefaultTbClusterServiceTest { |
||||
|
|
||||
|
public static final String MONOLITH = "monolith"; |
||||
|
|
||||
|
public static final String CORE = "core"; |
||||
|
|
||||
|
public static final String RULE_ENGINE = "rule_engine"; |
||||
|
|
||||
|
public static final String TRANSPORT = "transport"; |
||||
|
|
||||
|
@MockBean |
||||
|
protected DataDecodingEncodingService encodingService; |
||||
|
@MockBean |
||||
|
protected TbDeviceProfileCache deviceProfileCache; |
||||
|
@MockBean |
||||
|
protected TbAssetProfileCache assetProfileCache; |
||||
|
@MockBean |
||||
|
protected GatewayNotificationsService gatewayNotificationsService; |
||||
|
@MockBean |
||||
|
protected PartitionService partitionService; |
||||
|
@MockBean |
||||
|
protected TbQueueProducerProvider producerProvider; |
||||
|
|
||||
|
@SpyBean |
||||
|
protected NotificationsTopicService notificationsTopicService; |
||||
|
@SpyBean |
||||
|
protected TbClusterService clusterService; |
||||
|
|
||||
|
@Test |
||||
|
public void testOnQueueChangeSingleMonolith() { |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet(MONOLITH)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(MONOLITH)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(MONOLITH)); |
||||
|
|
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToRuleEngineNotificationMsg>> tbQueueProducer = mock(TbQueueProducer.class); |
||||
|
|
||||
|
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbQueueProducer); |
||||
|
|
||||
|
clusterService.onQueueChange(createTestQueue()); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(eq(ServiceType.TB_CORE), any()); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(eq(ServiceType.TB_TRANSPORT), any()); |
||||
|
|
||||
|
verify(tbQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(producerProvider, never()).getTbCoreNotificationsMsgProducer(); |
||||
|
verify(producerProvider, never()).getTransportNotificationsMsgProducer(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnQueueChangeMultipleMonoliths() { |
||||
|
String monolith1 = MONOLITH + 1; |
||||
|
String monolith2 = MONOLITH + 2; |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet(monolith1, monolith2)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(monolith1, monolith2)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(monolith1, monolith2)); |
||||
|
|
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToRuleEngineNotificationMsg>> tbQueueProducer = mock(TbQueueProducer.class); |
||||
|
|
||||
|
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbQueueProducer); |
||||
|
|
||||
|
clusterService.onQueueChange(createTestQueue()); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(eq(ServiceType.TB_CORE), any()); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(eq(ServiceType.TB_TRANSPORT), any()); |
||||
|
|
||||
|
verify(tbQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(producerProvider, never()).getTbCoreNotificationsMsgProducer(); |
||||
|
verify(producerProvider, never()).getTransportNotificationsMsgProducer(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnQueueChangeSingleMonolithAndSingleRemoteTransport() { |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet(MONOLITH)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(MONOLITH)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(MONOLITH, TRANSPORT)); |
||||
|
|
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToRuleEngineNotificationMsg>> tbREQueueProducer = mock(TbQueueProducer.class); |
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToTransportMsg>> tbTransportQueueProducer = mock(TbQueueProducer.class); |
||||
|
|
||||
|
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbREQueueProducer); |
||||
|
when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer); |
||||
|
|
||||
|
clusterService.onQueueChange(createTestQueue()); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_TRANSPORT, TRANSPORT); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(eq(ServiceType.TB_CORE), any()); |
||||
|
|
||||
|
verify(tbREQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(tbTransportQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, TRANSPORT)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbTransportQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, MONOLITH)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(tbTransportQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, MONOLITH)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(producerProvider, never()).getTbCoreNotificationsMsgProducer(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnQueueChangeMultipleMicroservices() { |
||||
|
String monolith1 = MONOLITH + 1; |
||||
|
String monolith2 = MONOLITH + 2; |
||||
|
|
||||
|
String core1 = CORE + 1; |
||||
|
String core2 = CORE + 2; |
||||
|
|
||||
|
String ruleEngine1 = RULE_ENGINE + 1; |
||||
|
String ruleEngine2 = RULE_ENGINE + 2; |
||||
|
|
||||
|
String transport1 = TRANSPORT + 1; |
||||
|
String transport2 = TRANSPORT + 2; |
||||
|
|
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet(monolith1, monolith2, ruleEngine1, ruleEngine2)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(monolith1, monolith2, core1, core2)); |
||||
|
when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(monolith1, monolith2, transport1, transport2)); |
||||
|
|
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToRuleEngineNotificationMsg>> tbREQueueProducer = mock(TbQueueProducer.class); |
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToCoreNotificationMsg>> tbCoreQueueProducer = mock(TbQueueProducer.class); |
||||
|
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToTransportMsg>> tbTransportQueueProducer = mock(TbQueueProducer.class); |
||||
|
|
||||
|
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbREQueueProducer); |
||||
|
when(producerProvider.getTbCoreNotificationsMsgProducer()).thenReturn(tbCoreQueueProducer); |
||||
|
when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer); |
||||
|
|
||||
|
clusterService.onQueueChange(createTestQueue()); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, ruleEngine1); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, ruleEngine2); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_CORE, core1); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_CORE, core2); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(ServiceType.TB_CORE, monolith1); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(ServiceType.TB_CORE, monolith2); |
||||
|
|
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_TRANSPORT, transport1); |
||||
|
verify(notificationsTopicService, times(1)).getNotificationsTopic(ServiceType.TB_TRANSPORT, transport2); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(ServiceType.TB_TRANSPORT, monolith1); |
||||
|
verify(notificationsTopicService, never()).getNotificationsTopic(ServiceType.TB_TRANSPORT, monolith2); |
||||
|
|
||||
|
verify(tbREQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbREQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbREQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, ruleEngine1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbREQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, ruleEngine2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(tbCoreQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, core1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbCoreQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, core2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbCoreQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, monolith1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbCoreQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, monolith2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
|
||||
|
verify(tbTransportQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, transport1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbTransportQueueProducer, times(1)) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, transport2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbTransportQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, monolith1)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
verify(tbTransportQueueProducer, never()) |
||||
|
.send(eq(notificationsTopicService.getNotificationsTopic(ServiceType.TB_TRANSPORT, monolith2)), any(TbProtoQueueMsg.class), isNull()); |
||||
|
} |
||||
|
|
||||
|
protected Queue createTestQueue() { |
||||
|
TenantId tenantId = TenantId.SYS_TENANT_ID; |
||||
|
Queue queue = new Queue(new QueueId(UUID.randomUUID())); |
||||
|
queue.setTenantId(tenantId); |
||||
|
queue.setName("Main"); |
||||
|
queue.setTopic("main"); |
||||
|
queue.setPartitions(10); |
||||
|
return queue; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,68 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.security.auth.oauth2; |
||||
|
|
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.mockito.Mock; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.thingsboard.server.common.data.id.UserId; |
||||
|
import org.thingsboard.server.common.data.security.model.JwtPair; |
||||
|
import org.thingsboard.server.controller.AbstractControllerTest; |
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
import org.thingsboard.server.service.security.model.SecurityUser; |
||||
|
import org.thingsboard.server.service.security.model.token.JwtTokenFactory; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.junit.Assert.assertEquals; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@DaoSqlTest |
||||
|
public class Oauth2AuthenticationSuccessHandlerTest extends AbstractControllerTest { |
||||
|
|
||||
|
@Autowired |
||||
|
private Oauth2AuthenticationSuccessHandler oauth2AuthenticationSuccessHandler; |
||||
|
|
||||
|
@Mock |
||||
|
private JwtTokenFactory jwtTokenFactory; |
||||
|
|
||||
|
private SecurityUser securityUser; |
||||
|
|
||||
|
@Before |
||||
|
public void before() { |
||||
|
UserId userId = new UserId(UUID.randomUUID()); |
||||
|
securityUser = new SecurityUser(userId); |
||||
|
when(jwtTokenFactory.createTokenPair(eq(securityUser))).thenReturn(new JwtPair("testAccessToken", "testRefreshToken")); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testGetRedirectUrl() { |
||||
|
JwtPair jwtPair = jwtTokenFactory.createTokenPair(securityUser); |
||||
|
|
||||
|
String urlWithoutParams = "http://localhost:8080/dashboardGroups/3fa13530-6597-11ed-bd76-8bd591f0ec3e"; |
||||
|
String urlWithParams = "http://localhost:8080/dashboardGroups/3fa13530-6597-11ed-bd76-8bd591f0ec3e?state=someState&page=1"; |
||||
|
|
||||
|
String redirectUrl = oauth2AuthenticationSuccessHandler.getRedirectUrl(urlWithoutParams, jwtPair); |
||||
|
String expectedUrl = urlWithoutParams + "/?accessToken=" + jwtPair.getToken() + "&refreshToken=" + jwtPair.getRefreshToken(); |
||||
|
assertEquals(expectedUrl, redirectUrl); |
||||
|
|
||||
|
redirectUrl = oauth2AuthenticationSuccessHandler.getRedirectUrl(urlWithParams, jwtPair); |
||||
|
expectedUrl = urlWithParams + "&accessToken=" + jwtPair.getToken() + "&refreshToken=" + jwtPair.getRefreshToken(); |
||||
|
assertEquals(expectedUrl, redirectUrl); |
||||
|
} |
||||
|
} |
||||
@ -1,167 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.script.api.mvel; |
|
||||
|
|
||||
import org.mvel2.ExecutionContext; |
|
||||
import org.mvel2.ParserConfiguration; |
|
||||
import org.mvel2.execution.ExecutionArrayList; |
|
||||
import org.mvel2.util.MethodStub; |
|
||||
|
|
||||
import java.io.UnsupportedEncodingException; |
|
||||
import java.util.Base64; |
|
||||
import java.util.List; |
|
||||
|
|
||||
public class TbUtils { |
|
||||
|
|
||||
public static void register(ParserConfiguration parserConfig) throws Exception { |
|
||||
parserConfig.addImport("btoa", new MethodStub(TbUtils.class.getMethod("btoa", |
|
||||
String.class))); |
|
||||
parserConfig.addImport("atob", new MethodStub(TbUtils.class.getMethod("atob", |
|
||||
String.class))); |
|
||||
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
|
||||
List.class))); |
|
||||
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
|
||||
List.class, String.class))); |
|
||||
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
|
||||
ExecutionContext.class, String.class))); |
|
||||
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
|
||||
ExecutionContext.class, String.class, String.class))); |
|
||||
parserConfig.addImport("parseInt", new MethodStub(TbUtils.class.getMethod("parseInt", |
|
||||
String.class))); |
|
||||
parserConfig.addImport("parseInt", new MethodStub(TbUtils.class.getMethod("parseInt", |
|
||||
String.class, int.class))); |
|
||||
parserConfig.addImport("parseFloat", new MethodStub(TbUtils.class.getMethod("parseFloat", |
|
||||
String.class))); |
|
||||
parserConfig.addImport("parseDouble", new MethodStub(TbUtils.class.getMethod("parseDouble", |
|
||||
String.class))); |
|
||||
} |
|
||||
|
|
||||
public static void main(String[] args) { |
|
||||
System.out.println(Integer.class == int.class); |
|
||||
} |
|
||||
|
|
||||
public static String btoa(String input) { |
|
||||
return new String(Base64.getEncoder().encode(input.getBytes())); |
|
||||
} |
|
||||
|
|
||||
public static String atob(String encoded) { |
|
||||
return new String(Base64.getDecoder().decode(encoded)); |
|
||||
} |
|
||||
|
|
||||
public static String bytesToString(List<Byte> bytesList) { |
|
||||
byte[] bytes = bytesFromList(bytesList); |
|
||||
return new String(bytes); |
|
||||
} |
|
||||
|
|
||||
public static String bytesToString(List<Byte> bytesList, String charsetName) throws UnsupportedEncodingException { |
|
||||
byte[] bytes = bytesFromList(bytesList); |
|
||||
return new String(bytes, charsetName); |
|
||||
} |
|
||||
|
|
||||
public static List<Byte> stringToBytes(ExecutionContext ctx, String str) { |
|
||||
byte[] bytes = str.getBytes(); |
|
||||
return bytesToList(ctx, bytes); |
|
||||
} |
|
||||
|
|
||||
public static List<Byte> stringToBytes(ExecutionContext ctx, String str, String charsetName) throws UnsupportedEncodingException { |
|
||||
byte[] bytes = str.getBytes(charsetName); |
|
||||
return bytesToList(ctx, bytes); |
|
||||
} |
|
||||
|
|
||||
private static byte[] bytesFromList(List<Byte> bytesList) { |
|
||||
byte[] bytes = new byte[bytesList.size()]; |
|
||||
for (int i = 0; i < bytesList.size(); i++) { |
|
||||
bytes[i] = bytesList.get(i); |
|
||||
} |
|
||||
return bytes; |
|
||||
} |
|
||||
|
|
||||
private static List<Byte> bytesToList(ExecutionContext ctx, byte[] bytes) { |
|
||||
List<Byte> list = new ExecutionArrayList<>(ctx); |
|
||||
for (int i = 0; i < bytes.length; i++) { |
|
||||
list.add(bytes[i]); |
|
||||
} |
|
||||
return list; |
|
||||
} |
|
||||
|
|
||||
public static Integer parseInt(String value) { |
|
||||
if (value != null) { |
|
||||
try { |
|
||||
int radix = 10; |
|
||||
if (isHexadecimal(value)) { |
|
||||
radix = 16; |
|
||||
} |
|
||||
return Integer.parseInt(prepareNumberString(value), radix); |
|
||||
} catch (NumberFormatException e) { |
|
||||
Float f = parseFloat(value); |
|
||||
if (f != null) { |
|
||||
return f.intValue(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
|
|
||||
public static Integer parseInt(String value, int radix) { |
|
||||
if (value != null) { |
|
||||
try { |
|
||||
return Integer.parseInt(prepareNumberString(value), radix); |
|
||||
} catch (NumberFormatException e) { |
|
||||
Float f = parseFloat(value); |
|
||||
if (f != null) { |
|
||||
return f.intValue(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
|
|
||||
public static Float parseFloat(String value) { |
|
||||
if (value != null) { |
|
||||
try { |
|
||||
return Float.parseFloat(prepareNumberString(value)); |
|
||||
} catch (NumberFormatException e) { |
|
||||
} |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
|
|
||||
public static Double parseDouble(String value) { |
|
||||
if (value != null) { |
|
||||
try { |
|
||||
return Double.parseDouble(prepareNumberString(value)); |
|
||||
} catch (NumberFormatException e) { |
|
||||
} |
|
||||
} |
|
||||
return null; |
|
||||
} |
|
||||
|
|
||||
private static boolean isHexadecimal(String value) { |
|
||||
return value != null && (value.contains("0x") || value.contains("0X")); |
|
||||
} |
|
||||
|
|
||||
private static String prepareNumberString(String value) { |
|
||||
if (value != null) { |
|
||||
value = value.trim(); |
|
||||
if (isHexadecimal(value)) { |
|
||||
value = value.replace("0x", ""); |
|
||||
value = value.replace("0X", ""); |
|
||||
} |
|
||||
value = value.replace(",", "."); |
|
||||
} |
|
||||
return value; |
|
||||
} |
|
||||
} |
|
||||
@ -0,0 +1,319 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.tbel; |
||||
|
|
||||
|
import org.mvel2.ExecutionContext; |
||||
|
import org.mvel2.ParserConfiguration; |
||||
|
import org.mvel2.execution.ExecutionArrayList; |
||||
|
import org.mvel2.util.MethodStub; |
||||
|
|
||||
|
import java.io.IOException; |
||||
|
import java.io.UnsupportedEncodingException; |
||||
|
import java.math.BigDecimal; |
||||
|
import java.math.RoundingMode; |
||||
|
import java.nio.ByteBuffer; |
||||
|
import java.nio.ByteOrder; |
||||
|
import java.nio.charset.StandardCharsets; |
||||
|
import java.util.Base64; |
||||
|
import java.util.List; |
||||
|
|
||||
|
public class TbUtils { |
||||
|
|
||||
|
private static final byte[] HEX_ARRAY = "0123456789ABCDEF".getBytes(StandardCharsets.US_ASCII); |
||||
|
|
||||
|
public static void register(ParserConfiguration parserConfig) throws Exception { |
||||
|
parserConfig.addImport("btoa", new MethodStub(TbUtils.class.getMethod("btoa", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("atob", new MethodStub(TbUtils.class.getMethod("atob", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
||||
|
List.class))); |
||||
|
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
||||
|
List.class, String.class))); |
||||
|
parserConfig.addImport("decodeToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
||||
|
List.class))); |
||||
|
parserConfig.addImport("decodeToJson", new MethodStub(TbUtils.class.getMethod("decodeToJson", |
||||
|
ExecutionContext.class, List.class))); |
||||
|
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
||||
|
ExecutionContext.class, String.class))); |
||||
|
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
||||
|
ExecutionContext.class, String.class, String.class))); |
||||
|
parserConfig.addImport("parseInt", new MethodStub(TbUtils.class.getMethod("parseInt", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseInt", new MethodStub(TbUtils.class.getMethod("parseInt", |
||||
|
String.class, int.class))); |
||||
|
parserConfig.addImport("parseFloat", new MethodStub(TbUtils.class.getMethod("parseFloat", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseDouble", new MethodStub(TbUtils.class.getMethod("parseDouble", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseLittleEndianHexToInt", new MethodStub(TbUtils.class.getMethod("parseLittleEndianHexToInt", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseBigEndianHexToInt", new MethodStub(TbUtils.class.getMethod("parseBigEndianHexToInt", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseHexToInt", new MethodStub(TbUtils.class.getMethod("parseHexToInt", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("parseHexToInt", new MethodStub(TbUtils.class.getMethod("parseHexToInt", |
||||
|
String.class, boolean.class))); |
||||
|
parserConfig.addImport("parseBytesToInt", new MethodStub(TbUtils.class.getMethod("parseBytesToInt", |
||||
|
List.class, int.class, int.class))); |
||||
|
parserConfig.addImport("parseBytesToInt", new MethodStub(TbUtils.class.getMethod("parseBytesToInt", |
||||
|
List.class, int.class, int.class, boolean.class))); |
||||
|
parserConfig.addImport("parseBytesToInt", new MethodStub(TbUtils.class.getMethod("parseBytesToInt", |
||||
|
byte[].class, int.class, int.class))); |
||||
|
parserConfig.addImport("parseBytesToInt", new MethodStub(TbUtils.class.getMethod("parseBytesToInt", |
||||
|
byte[].class, int.class, int.class, boolean.class))); |
||||
|
parserConfig.addImport("toFixed", new MethodStub(TbUtils.class.getMethod("toFixed", |
||||
|
double.class, int.class))); |
||||
|
parserConfig.addImport("hexToBytes", new MethodStub(TbUtils.class.getMethod("hexToBytes", |
||||
|
ExecutionContext.class, String.class))); |
||||
|
parserConfig.addImport("base64ToHex", new MethodStub(TbUtils.class.getMethod("base64ToHex", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("base64ToBytes", new MethodStub(TbUtils.class.getMethod("base64ToBytes", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("bytesToBase64", new MethodStub(TbUtils.class.getMethod("bytesToBase64", |
||||
|
byte[].class))); |
||||
|
parserConfig.addImport("bytesToHex", new MethodStub(TbUtils.class.getMethod("bytesToHex", |
||||
|
byte[].class))); |
||||
|
parserConfig.addImport("bytesToHex", new MethodStub(TbUtils.class.getMethod("bytesToHex", |
||||
|
ExecutionArrayList.class))); |
||||
|
} |
||||
|
|
||||
|
public static String btoa(String input) { |
||||
|
return new String(Base64.getEncoder().encode(input.getBytes())); |
||||
|
} |
||||
|
|
||||
|
public static String atob(String encoded) { |
||||
|
return new String(Base64.getDecoder().decode(encoded)); |
||||
|
} |
||||
|
|
||||
|
public static Object decodeToJson(ExecutionContext ctx, List<Byte> bytesList) throws IOException { |
||||
|
return TbJson.parse(ctx, bytesToString(bytesList)); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToString(List<Byte> bytesList) { |
||||
|
byte[] bytes = bytesFromList(bytesList); |
||||
|
return new String(bytes); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToString(List<Byte> bytesList, String charsetName) throws UnsupportedEncodingException { |
||||
|
byte[] bytes = bytesFromList(bytesList); |
||||
|
return new String(bytes, charsetName); |
||||
|
} |
||||
|
|
||||
|
public static List<Byte> stringToBytes(ExecutionContext ctx, String str) { |
||||
|
byte[] bytes = str.getBytes(); |
||||
|
return bytesToList(ctx, bytes); |
||||
|
} |
||||
|
|
||||
|
public static List<Byte> stringToBytes(ExecutionContext ctx, String str, String charsetName) throws UnsupportedEncodingException { |
||||
|
byte[] bytes = str.getBytes(charsetName); |
||||
|
return bytesToList(ctx, bytes); |
||||
|
} |
||||
|
|
||||
|
private static byte[] bytesFromList(List<Byte> bytesList) { |
||||
|
byte[] bytes = new byte[bytesList.size()]; |
||||
|
for (int i = 0; i < bytesList.size(); i++) { |
||||
|
bytes[i] = bytesList.get(i); |
||||
|
} |
||||
|
return bytes; |
||||
|
} |
||||
|
|
||||
|
private static List<Byte> bytesToList(ExecutionContext ctx, byte[] bytes) { |
||||
|
List<Byte> list = new ExecutionArrayList<>(ctx); |
||||
|
for (byte aByte : bytes) { |
||||
|
list.add(aByte); |
||||
|
} |
||||
|
return list; |
||||
|
} |
||||
|
|
||||
|
public static Integer parseInt(String value) { |
||||
|
if (value != null) { |
||||
|
try { |
||||
|
int radix = 10; |
||||
|
if (isHexadecimal(value)) { |
||||
|
radix = 16; |
||||
|
} |
||||
|
return Integer.parseInt(prepareNumberString(value), radix); |
||||
|
} catch (NumberFormatException e) { |
||||
|
Float f = parseFloat(value); |
||||
|
if (f != null) { |
||||
|
return f.intValue(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
public static Integer parseInt(String value, int radix) { |
||||
|
if (value != null) { |
||||
|
try { |
||||
|
return Integer.parseInt(prepareNumberString(value), radix); |
||||
|
} catch (NumberFormatException e) { |
||||
|
Float f = parseFloat(value); |
||||
|
if (f != null) { |
||||
|
return f.intValue(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
public static Float parseFloat(String value) { |
||||
|
if (value != null) { |
||||
|
try { |
||||
|
return Float.parseFloat(prepareNumberString(value)); |
||||
|
} catch (NumberFormatException e) { |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
public static Double parseDouble(String value) { |
||||
|
if (value != null) { |
||||
|
try { |
||||
|
return Double.parseDouble(prepareNumberString(value)); |
||||
|
} catch (NumberFormatException e) { |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
public static int parseLittleEndianHexToInt(String hex) { |
||||
|
return parseHexToInt(hex, false); |
||||
|
} |
||||
|
|
||||
|
public static int parseBigEndianHexToInt(String hex) { |
||||
|
return parseHexToInt(hex, true); |
||||
|
} |
||||
|
|
||||
|
public static int parseHexToInt(String hex) { |
||||
|
return parseHexToInt(hex, true); |
||||
|
} |
||||
|
|
||||
|
public static int parseHexToInt(String hex, boolean bigEndian) { |
||||
|
int length = hex.length(); |
||||
|
if (length > 8) { |
||||
|
throw new IllegalArgumentException("Hex string is too large. Maximum 8 symbols allowed."); |
||||
|
} |
||||
|
if (length % 2 > 0) { |
||||
|
throw new IllegalArgumentException("Hex string must be even-length."); |
||||
|
} |
||||
|
byte[] data = new byte[length / 2]; |
||||
|
for (int i = 0; i < length; i += 2) { |
||||
|
data[i / 2] = (byte) ((Character.digit(hex.charAt(i), 16) << 4) + Character.digit(hex.charAt(i + 1), 16)); |
||||
|
} |
||||
|
return parseBytesToInt(data, 0, data.length, bigEndian); |
||||
|
} |
||||
|
|
||||
|
public static ExecutionArrayList<Byte> hexToBytes(ExecutionContext ctx, String hex) { |
||||
|
int len = hex.length(); |
||||
|
if (len % 2 > 0) { |
||||
|
throw new IllegalArgumentException("Hex string must be even-length."); |
||||
|
} |
||||
|
ExecutionArrayList<Byte> data = new ExecutionArrayList<>(ctx); |
||||
|
for (int i = 0; i < len; i += 2) { |
||||
|
data.add((byte)((Character.digit(hex.charAt(i), 16) << 4) |
||||
|
+ Character.digit(hex.charAt(i + 1), 16))); |
||||
|
} |
||||
|
return data; |
||||
|
} |
||||
|
|
||||
|
public static String base64ToHex(String base64) { |
||||
|
return bytesToHex(Base64.getDecoder().decode(base64)); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToBase64(byte[] bytes) { |
||||
|
return Base64.getEncoder().encodeToString(bytes); |
||||
|
} |
||||
|
|
||||
|
public static byte[] base64ToBytes(String input) { |
||||
|
return Base64.getDecoder().decode(input); |
||||
|
} |
||||
|
|
||||
|
public static int parseBytesToInt(List<Byte> data, int offset, int length) { |
||||
|
return parseBytesToInt(data, offset, length, true); |
||||
|
} |
||||
|
|
||||
|
public static int parseBytesToInt(List<Byte> data, int offset, int length, boolean bigEndian) { |
||||
|
final byte[] bytes = new byte[data.size()]; |
||||
|
for (int i = 0; i < bytes.length; i++) { |
||||
|
bytes[i] = data.get(i); |
||||
|
} |
||||
|
return parseBytesToInt(bytes, offset, length, bigEndian); |
||||
|
} |
||||
|
|
||||
|
public static int parseBytesToInt(byte[] data, int offset, int length) { |
||||
|
return parseBytesToInt(data, offset, length, true); |
||||
|
} |
||||
|
|
||||
|
public static int parseBytesToInt(byte[] data, int offset, int length, boolean bigEndian) { |
||||
|
if (offset > data.length) { |
||||
|
throw new IllegalArgumentException("Offset: " + offset + " is out of bounds for array with length: " + data.length + "!"); |
||||
|
} |
||||
|
if (length > 4) { |
||||
|
throw new IllegalArgumentException("Length: " + length + " is too large. Maximum 4 bytes is allowed!"); |
||||
|
} |
||||
|
if (offset + length > data.length) { |
||||
|
throw new IllegalArgumentException("Offset: " + offset + " and Length: " + length + " is out of bounds for array with length: " + data.length + "!"); |
||||
|
} |
||||
|
var bb = ByteBuffer.allocate(4); |
||||
|
if (!bigEndian) { |
||||
|
bb.order(ByteOrder.LITTLE_ENDIAN); |
||||
|
} |
||||
|
bb.position(bigEndian ? 4 - length : 0); |
||||
|
bb.put(data, offset, length); |
||||
|
bb.position(0); |
||||
|
return bb.getInt(); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToHex(ExecutionArrayList<?> bytesList) { |
||||
|
byte[] bytes = new byte[bytesList.size()]; |
||||
|
for (int i = 0; i < bytesList.size(); i++) { |
||||
|
bytes[i] = Byte.parseByte(bytesList.get(i).toString()); |
||||
|
} |
||||
|
return bytesToHex(bytes); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToHex(byte[] bytes) { |
||||
|
byte[] hexChars = new byte[bytes.length * 2]; |
||||
|
for (int j = 0; j < bytes.length; j++) { |
||||
|
int v = bytes[j] & 0xFF; |
||||
|
hexChars[j * 2] = HEX_ARRAY[v >>> 4]; |
||||
|
hexChars[j * 2 + 1] = HEX_ARRAY[v & 0x0F]; |
||||
|
} |
||||
|
return new String(hexChars, StandardCharsets.UTF_8); |
||||
|
} |
||||
|
|
||||
|
public static double toFixed(double value, int precision) { |
||||
|
return BigDecimal.valueOf(value).setScale(precision, RoundingMode.HALF_UP).doubleValue(); |
||||
|
} |
||||
|
|
||||
|
private static boolean isHexadecimal(String value) { |
||||
|
return value != null && (value.contains("0x") || value.contains("0X")); |
||||
|
} |
||||
|
|
||||
|
private static String prepareNumberString(String value) { |
||||
|
if (value != null) { |
||||
|
value = value.trim(); |
||||
|
if (isHexadecimal(value)) { |
||||
|
value = value.replace("0x", ""); |
||||
|
value = value.replace("0X", ""); |
||||
|
} |
||||
|
value = value.replace(",", "."); |
||||
|
} |
||||
|
return value; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,98 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.tbel; |
||||
|
|
||||
|
import org.junit.Assert; |
||||
|
import org.junit.Test; |
||||
|
|
||||
|
import java.nio.ByteBuffer; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
|
||||
|
public class TbUtilsTest { |
||||
|
|
||||
|
@Test |
||||
|
public void parseHexToInt() { |
||||
|
Assert.assertEquals(0xAB, TbUtils.parseHexToInt("AB")); |
||||
|
Assert.assertEquals(0xABBA, TbUtils.parseHexToInt("ABBA", true)); |
||||
|
Assert.assertEquals(0xBAAB, TbUtils.parseHexToInt("ABBA", false)); |
||||
|
Assert.assertEquals(0xAABBCC, TbUtils.parseHexToInt("AABBCC", true)); |
||||
|
Assert.assertEquals(0xAABBCC, TbUtils.parseHexToInt("CCBBAA", false)); |
||||
|
Assert.assertEquals(0xAABBCCDD, TbUtils.parseHexToInt("AABBCCDD", true)); |
||||
|
Assert.assertEquals(0xAABBCCDD, TbUtils.parseHexToInt("DDCCBBAA", false)); |
||||
|
Assert.assertEquals(0xDDCCBBAA, TbUtils.parseHexToInt("DDCCBBAA", true)); |
||||
|
Assert.assertEquals(0xDDCCBBAA, TbUtils.parseHexToInt("AABBCCDD", false)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void parseBytesToInt_checkPrimitives() { |
||||
|
int expected = 257; |
||||
|
byte[] data = ByteBuffer.allocate(4).putInt(expected).array(); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4)); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 2, 2, true)); |
||||
|
Assert.assertEquals(1, TbUtils.parseBytesToInt(data, 3, 1, true)); |
||||
|
|
||||
|
expected = Integer.MAX_VALUE; |
||||
|
data = ByteBuffer.allocate(4).putInt(expected).array(); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, true)); |
||||
|
|
||||
|
expected = 0xAABBCCDD; |
||||
|
data = new byte[]{(byte) 0xAA, (byte) 0xBB, (byte) 0xCC, (byte) 0xDD}; |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, true)); |
||||
|
data = new byte[]{(byte) 0xDD, (byte) 0xCC, (byte) 0xBB, (byte) 0xAA}; |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, false)); |
||||
|
|
||||
|
expected = 0xAABBCC; |
||||
|
data = new byte[]{(byte) 0xAA, (byte) 0xBB, (byte) 0xCC}; |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 3, true)); |
||||
|
data = new byte[]{(byte) 0xCC, (byte) 0xBB, (byte) 0xAA}; |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 3, false)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void parseBytesToInt_checkLists() { |
||||
|
int expected = 257; |
||||
|
List<Byte> data = toList(ByteBuffer.allocate(4).putInt(expected).array()); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4)); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 2, 2, true)); |
||||
|
Assert.assertEquals(1, TbUtils.parseBytesToInt(data, 3, 1, true)); |
||||
|
|
||||
|
expected = Integer.MAX_VALUE; |
||||
|
data = toList(ByteBuffer.allocate(4).putInt(expected).array()); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, true)); |
||||
|
|
||||
|
expected = 0xAABBCCDD; |
||||
|
data = toList(new byte[]{(byte) 0xAA, (byte) 0xBB, (byte) 0xCC, (byte) 0xDD}); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, true)); |
||||
|
data = toList(new byte[]{(byte) 0xDD, (byte) 0xCC, (byte) 0xBB, (byte) 0xAA}); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 4, false)); |
||||
|
|
||||
|
expected = 0xAABBCC; |
||||
|
data = toList(new byte[]{(byte) 0xAA, (byte) 0xBB, (byte) 0xCC}); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 3, true)); |
||||
|
data = toList(new byte[]{(byte) 0xCC, (byte) 0xBB, (byte) 0xAA}); |
||||
|
Assert.assertEquals(expected, TbUtils.parseBytesToInt(data, 0, 3, false)); |
||||
|
} |
||||
|
|
||||
|
private static List<Byte> toList(byte[] data) { |
||||
|
List<Byte> result = new ArrayList<>(data.length); |
||||
|
for (Byte b : data) { |
||||
|
result.add(b); |
||||
|
} |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,121 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.msa; |
||||
|
|
||||
|
import java.io.IOException; |
||||
|
import org.eclipse.californium.core.CoapClient; |
||||
|
import org.eclipse.californium.core.CoapHandler; |
||||
|
import org.eclipse.californium.core.CoapObserveRelation; |
||||
|
import org.eclipse.californium.core.CoapResponse; |
||||
|
import org.eclipse.californium.core.coap.CoAP; |
||||
|
import org.eclipse.californium.core.coap.MediaTypeRegistry; |
||||
|
import org.eclipse.californium.core.coap.Request; |
||||
|
import org.eclipse.californium.elements.exception.ConnectorException; |
||||
|
import org.thingsboard.server.common.msg.session.FeatureType; |
||||
|
|
||||
|
public class TestCoapClient { |
||||
|
|
||||
|
private static final String COAP_BASE_URL = "coap://localhost:5683/api/v1/"; |
||||
|
private static final long CLIENT_REQUEST_TIMEOUT = 60000L; |
||||
|
|
||||
|
private final CoapClient client; |
||||
|
|
||||
|
public TestCoapClient(){ |
||||
|
this.client = createClient(); |
||||
|
} |
||||
|
|
||||
|
public TestCoapClient(String accessToken, FeatureType featureType) { |
||||
|
this.client = createClient(getFeatureTokenUrl(accessToken, featureType)); |
||||
|
} |
||||
|
|
||||
|
public TestCoapClient(String featureTokenUrl) { |
||||
|
this.client = createClient(featureTokenUrl); |
||||
|
} |
||||
|
|
||||
|
public void connectToCoap(String accessToken) { |
||||
|
setURI(accessToken, null); |
||||
|
} |
||||
|
|
||||
|
public void connectToCoap(String accessToken, FeatureType featureType) { |
||||
|
setURI(accessToken, featureType); |
||||
|
} |
||||
|
|
||||
|
public void disconnect() { |
||||
|
if (client != null) { |
||||
|
client.shutdown(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public CoapResponse postMethod(String requestBody) throws ConnectorException, IOException { |
||||
|
return this.postMethod(requestBody.getBytes()); |
||||
|
} |
||||
|
|
||||
|
public CoapResponse postMethod(byte[] requestBodyBytes) throws ConnectorException, IOException { |
||||
|
return client.setTimeout(CLIENT_REQUEST_TIMEOUT).post(requestBodyBytes, MediaTypeRegistry.APPLICATION_JSON); |
||||
|
} |
||||
|
|
||||
|
public void postMethod(CoapHandler handler, String payload, int format) { |
||||
|
client.post(handler, payload, format); |
||||
|
} |
||||
|
|
||||
|
public void postMethod(CoapHandler handler, byte[] payload, int format) { |
||||
|
client.post(handler, payload, format); |
||||
|
} |
||||
|
|
||||
|
public CoapResponse getMethod() throws ConnectorException, IOException { |
||||
|
return client.setTimeout(CLIENT_REQUEST_TIMEOUT).get(); |
||||
|
} |
||||
|
|
||||
|
public CoapObserveRelation getObserveRelation(TestCoapClientCallback callback){ |
||||
|
Request request = Request.newGet().setObserve(); |
||||
|
request.setType(CoAP.Type.CON); |
||||
|
return client.observe(request, callback); |
||||
|
} |
||||
|
|
||||
|
public void setURI(String featureTokenUrl) { |
||||
|
if (client == null) { |
||||
|
throw new RuntimeException("Failed to connect! CoapClient is not initialized!"); |
||||
|
} |
||||
|
client.setURI(featureTokenUrl); |
||||
|
} |
||||
|
|
||||
|
public void setURI(String accessToken, FeatureType featureType) { |
||||
|
if (featureType == null){ |
||||
|
featureType = FeatureType.ATTRIBUTES; |
||||
|
} |
||||
|
setURI(getFeatureTokenUrl(accessToken, featureType)); |
||||
|
} |
||||
|
|
||||
|
private CoapClient createClient() { |
||||
|
return new CoapClient(); |
||||
|
} |
||||
|
|
||||
|
private CoapClient createClient(String featureTokenUrl) { |
||||
|
return new CoapClient(featureTokenUrl); |
||||
|
} |
||||
|
|
||||
|
public static String getFeatureTokenUrl(FeatureType featureType) { |
||||
|
return COAP_BASE_URL + featureType.name().toLowerCase(); |
||||
|
} |
||||
|
|
||||
|
public static String getFeatureTokenUrl(String token, FeatureType featureType) { |
||||
|
return COAP_BASE_URL + token + "/" + featureType.name().toLowerCase(); |
||||
|
} |
||||
|
|
||||
|
public static String getFeatureTokenUrl(String token, FeatureType featureType, int requestId) { |
||||
|
return COAP_BASE_URL + token + "/" + featureType.name().toLowerCase() + "/" + requestId; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,68 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.msa; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.eclipse.californium.core.CoapHandler; |
||||
|
import org.eclipse.californium.core.CoapResponse; |
||||
|
import org.eclipse.californium.core.coap.CoAP; |
||||
|
|
||||
|
import java.util.concurrent.CountDownLatch; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Data |
||||
|
public class TestCoapClientCallback implements CoapHandler { |
||||
|
|
||||
|
protected final CountDownLatch latch; |
||||
|
protected Integer observe; |
||||
|
protected byte[] payloadBytes; |
||||
|
protected CoAP.ResponseCode responseCode; |
||||
|
|
||||
|
public TestCoapClientCallback() { |
||||
|
this.latch = new CountDownLatch(1); |
||||
|
} |
||||
|
|
||||
|
public TestCoapClientCallback(int subscribeCount) { |
||||
|
this.latch = new CountDownLatch(subscribeCount); |
||||
|
} |
||||
|
|
||||
|
public Integer getObserve() { |
||||
|
return observe; |
||||
|
} |
||||
|
|
||||
|
public byte[] getPayloadBytes() { |
||||
|
return payloadBytes; |
||||
|
} |
||||
|
|
||||
|
public CoAP.ResponseCode getResponseCode() { |
||||
|
return responseCode; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onLoad(CoapResponse response) { |
||||
|
observe = response.getOptions().getObserve(); |
||||
|
payloadBytes = response.getPayload(); |
||||
|
responseCode = response.getCode(); |
||||
|
latch.countDown(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onError() { |
||||
|
log.warn("Command Response Ack Error, No connect"); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,120 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.msa.connectivity; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import com.google.gson.JsonObject; |
||||
|
import io.restassured.path.json.JsonPath; |
||||
|
import org.testng.annotations.AfterMethod; |
||||
|
import org.testng.annotations.BeforeMethod; |
||||
|
import org.testng.annotations.Test; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
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.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.common.msg.session.FeatureType; |
||||
|
import org.thingsboard.server.msa.AbstractContainerTest; |
||||
|
import org.thingsboard.server.msa.TestCoapClient; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.thingsboard.server.msa.prototypes.DevicePrototypes.defaultDevicePrototype; |
||||
|
|
||||
|
public class CoapClientTest extends AbstractContainerTest { |
||||
|
private TestCoapClient client; |
||||
|
|
||||
|
private Device device; |
||||
|
@BeforeMethod |
||||
|
public void setUp() throws Exception { |
||||
|
testRestClient.login("tenant@thingsboard.org", "tenant"); |
||||
|
device = testRestClient.postDevice("", defaultDevicePrototype("http_")); |
||||
|
} |
||||
|
|
||||
|
@AfterMethod |
||||
|
public void tearDown() { |
||||
|
testRestClient.deleteDeviceIfExists(device.getId()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void provisionRequestForDeviceWithPreProvisionedStrategy() throws Exception { |
||||
|
|
||||
|
DeviceProfile deviceProfile = testRestClient.getDeviceProfileById(device.getDeviceProfileId()); |
||||
|
deviceProfile = updateDeviceProfileWithProvisioningStrategy(deviceProfile, DeviceProfileProvisionType.CHECK_PRE_PROVISIONED_DEVICES); |
||||
|
|
||||
|
DeviceCredentials expectedDeviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(device.getId()); |
||||
|
|
||||
|
JsonNode provisionResponse = JacksonUtil.fromBytes(createCoapClientAndPublish(device.getName())); |
||||
|
|
||||
|
assertThat(provisionResponse.get("credentialsType").asText()).isEqualTo(expectedDeviceCredentials.getCredentialsType().name()); |
||||
|
assertThat(provisionResponse.get("credentialsValue").asText()).isEqualTo(expectedDeviceCredentials.getCredentialsId()); |
||||
|
assertThat(provisionResponse.get("status").asText()).isEqualTo("SUCCESS"); |
||||
|
|
||||
|
updateDeviceProfileWithProvisioningStrategy(deviceProfile, DeviceProfileProvisionType.DISABLED); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void provisionRequestForDeviceWithAllowToCreateNewDevicesStrategy() throws Exception { |
||||
|
|
||||
|
String testDeviceName = "test_provision_device"; |
||||
|
|
||||
|
DeviceProfile deviceProfile = testRestClient.getDeviceProfileById(device.getDeviceProfileId()); |
||||
|
|
||||
|
deviceProfile = updateDeviceProfileWithProvisioningStrategy(deviceProfile, DeviceProfileProvisionType.ALLOW_CREATE_NEW_DEVICES); |
||||
|
|
||||
|
JsonNode provisionResponse = JacksonUtil.fromBytes(createCoapClientAndPublish(testDeviceName)); |
||||
|
|
||||
|
testRestClient.deleteDeviceIfExists(device.getId()); |
||||
|
device = testRestClient.getDeviceByName(testDeviceName); |
||||
|
|
||||
|
DeviceCredentials expectedDeviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(device.getId()); |
||||
|
|
||||
|
assertThat(provisionResponse.get("credentialsType").asText()).isEqualTo(expectedDeviceCredentials.getCredentialsType().name()); |
||||
|
assertThat(provisionResponse.get("credentialsValue").asText()).isEqualTo(expectedDeviceCredentials.getCredentialsId()); |
||||
|
assertThat(provisionResponse.get("status").asText()).isEqualTo("SUCCESS"); |
||||
|
|
||||
|
updateDeviceProfileWithProvisioningStrategy(deviceProfile, DeviceProfileProvisionType.DISABLED); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void provisionRequestForDeviceWithDisabledProvisioningStrategy() throws Exception { |
||||
|
|
||||
|
JsonObject provisionRequest = new JsonObject(); |
||||
|
provisionRequest.addProperty("provisionDeviceKey", TEST_PROVISION_DEVICE_KEY); |
||||
|
provisionRequest.addProperty("provisionDeviceSecret", TEST_PROVISION_DEVICE_SECRET); |
||||
|
|
||||
|
JsonNode response = JacksonUtil.fromBytes(createCoapClientAndPublish(null)); |
||||
|
|
||||
|
assertThat(response.get("status").asText()).isEqualTo("NOT_FOUND"); |
||||
|
} |
||||
|
|
||||
|
private byte[] createCoapClientAndPublish(String deviceName) throws Exception { |
||||
|
String provisionRequestMsg = createTestProvisionMessage(deviceName); |
||||
|
client = new TestCoapClient(TestCoapClient.getFeatureTokenUrl(FeatureType.PROVISION)); |
||||
|
return client.postMethod(provisionRequestMsg.getBytes()).getPayload(); |
||||
|
} |
||||
|
|
||||
|
private String createTestProvisionMessage(String deviceName) { |
||||
|
ObjectNode provisionRequest = JacksonUtil.newObjectNode(); |
||||
|
provisionRequest.put("provisionDeviceKey", TEST_PROVISION_DEVICE_KEY); |
||||
|
provisionRequest.put("provisionDeviceSecret", TEST_PROVISION_DEVICE_SECRET); |
||||
|
if (deviceName != null) { |
||||
|
provisionRequest.put("deviceName", deviceName); |
||||
|
} |
||||
|
return provisionRequest.toString(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue