432 changed files with 20009 additions and 4725 deletions
@ -0,0 +1,88 @@ |
|||
# |
|||
# 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. |
|||
# |
|||
|
|||
changelog: |
|||
exclude: |
|||
labels: |
|||
- Ignore for release |
|||
categories: |
|||
- title: 'Major Core & Rule Engine' |
|||
labels: |
|||
- 'Major Core' |
|||
- 'Major Rule Engine' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Major UI' |
|||
labels: |
|||
- 'Major UI' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Major Transport' |
|||
labels: |
|||
- 'Major Transport' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Major Edge' |
|||
labels: |
|||
- 'Major Edge' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Core & Rule Engine' |
|||
labels: |
|||
- 'Core' |
|||
- 'Rule Engine' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'UI' |
|||
labels: |
|||
- 'UI' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Transport' |
|||
labels: |
|||
- 'Transport' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Edge' |
|||
labels: |
|||
- 'Edge' |
|||
exclude: |
|||
labels: |
|||
- 'Bug' |
|||
- title: 'Bug: Core & Rule Engine' |
|||
labels: |
|||
- 'Core' |
|||
- 'Rule Engine' |
|||
- 'Bug' |
|||
- title: 'Bug: UI' |
|||
labels: |
|||
- 'UI' |
|||
- 'Bug' |
|||
- title: 'Bug: Transport' |
|||
labels: |
|||
- 'Transport' |
|||
- 'Bug' |
|||
- title: 'Bug: Edge' |
|||
labels: |
|||
- 'Edge' |
|||
- 'Bug' |
|||
@ -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; |
|||
} |
|||
} |
|||
4
application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java → application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
4
application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java → application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/request/MqttAttributesRequestBackwardCompatibilityIntegrationTest.java
4
application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java → application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
4
application/src/test/java/org/thingsboard/server/transport/mqtt/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java → application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/updates/MqttAttributesUpdatesBackwardCompatibilityIntegrationTest.java
@ -0,0 +1,50 @@ |
|||
/** |
|||
* 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.transport.mqtt.mqttv3.client; |
|||
|
|||
import org.eclipse.paho.client.mqttv3.MqttException; |
|||
import org.junit.Assert; |
|||
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; |
|||
import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; |
|||
|
|||
public abstract class AbstractMqttClientConnectionTest extends AbstractMqttIntegrationTest { |
|||
|
|||
protected void processClientWithCorrectAccessTokenTest() throws Exception { |
|||
MqttTestClient client = new MqttTestClient(); |
|||
client.connectAndWait(accessToken); |
|||
Assert.assertTrue(client.isConnected()); |
|||
client.disconnect(); |
|||
} |
|||
|
|||
protected void processClientWithWrongAccessTokenTest() throws Exception { |
|||
MqttTestClient client = new MqttTestClient(); |
|||
try { |
|||
client.connectAndWait("wrongAccessToken"); |
|||
} catch (MqttException e) { |
|||
Assert.assertEquals(MqttException.REASON_CODE_FAILED_AUTHENTICATION, e.getReasonCode()); |
|||
} |
|||
} |
|||
|
|||
protected void processClientWithWrongClientIdAndEmptyUsernamePasswordTest() throws Exception { |
|||
MqttTestClient client = new MqttTestClient("unknownClientId"); |
|||
try { |
|||
client.connectAndWait(); |
|||
} catch (MqttException e) { |
|||
Assert.assertEquals(MqttException.REASON_CODE_INVALID_CLIENT_ID, e.getReasonCode()); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,48 @@ |
|||
/** |
|||
* 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.transport.mqtt.mqttv3.client; |
|||
|
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; |
|||
|
|||
@DaoSqlTest |
|||
public class MqttClientConnectionTest extends AbstractMqttClientConnectionTest { |
|||
|
|||
@Before |
|||
public void beforeTest() throws Exception { |
|||
MqttTestConfigProperties configProperties = MqttTestConfigProperties.builder() |
|||
.deviceName("Test MqttV5 client device") |
|||
.build(); |
|||
processBeforeTest(configProperties); |
|||
} |
|||
|
|||
@Test |
|||
public void testClientWithCorrectAccessToken() throws Exception { |
|||
processClientWithCorrectAccessTokenTest(); |
|||
} |
|||
|
|||
@Test |
|||
public void testClientWithWrongAccessToken() throws Exception { |
|||
processClientWithWrongAccessTokenTest(); |
|||
} |
|||
|
|||
@Test |
|||
public void testClientWithWrongClientIdAndEmptyUsernamePassword() throws Exception { |
|||
processClientWithWrongClientIdAndEmptyUsernamePasswordTest(); |
|||
} |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
/** |
|||
* 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.transport.mqtt.mqttv5; |
|||
|
|||
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; |
|||
|
|||
public abstract class AbstractMqttV5Test extends AbstractMqttIntegrationTest { |
|||
} |
|||
@ -0,0 +1,110 @@ |
|||
/** |
|||
* 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.transport.mqtt.mqttv5; |
|||
|
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.paho.mqttv5.client.IMqttToken; |
|||
import org.eclipse.paho.mqttv5.client.MqttCallback; |
|||
import org.eclipse.paho.mqttv5.client.MqttDisconnectResponse; |
|||
import org.eclipse.paho.mqttv5.common.MqttException; |
|||
import org.eclipse.paho.mqttv5.common.MqttMessage; |
|||
import org.eclipse.paho.mqttv5.common.packet.MqttProperties; |
|||
import org.eclipse.paho.mqttv5.common.packet.MqttWireMessage; |
|||
|
|||
import java.util.concurrent.CountDownLatch; |
|||
|
|||
@Slf4j |
|||
@Data |
|||
public class MqttV5TestCallback implements MqttCallback { |
|||
|
|||
protected CountDownLatch subscribeLatch; |
|||
protected final CountDownLatch deliveryLatch; |
|||
protected int qoS; |
|||
protected byte[] payloadBytes; |
|||
protected String awaitSubTopic; |
|||
protected boolean pubAckReceived; |
|||
protected MqttMessage lastReceivedMessage; |
|||
|
|||
public MqttV5TestCallback() { |
|||
this.subscribeLatch = new CountDownLatch(1); |
|||
this.deliveryLatch = new CountDownLatch(1); |
|||
} |
|||
|
|||
public MqttV5TestCallback(int subscribeCount) { |
|||
this.subscribeLatch = new CountDownLatch(subscribeCount); |
|||
this.deliveryLatch = new CountDownLatch(1); |
|||
} |
|||
|
|||
public MqttV5TestCallback(String awaitSubTopic) { |
|||
this.subscribeLatch = new CountDownLatch(1); |
|||
this.deliveryLatch = new CountDownLatch(1); |
|||
this.awaitSubTopic = awaitSubTopic; |
|||
} |
|||
|
|||
@Override |
|||
public void disconnected(MqttDisconnectResponse mqttDisconnectResponse) { |
|||
if (mqttDisconnectResponse.getException() != null) { |
|||
log.warn("connectionLost: ", mqttDisconnectResponse.getException()); |
|||
deliveryLatch.countDown(); |
|||
} |
|||
log.warn("Disconnected with reason: {}", mqttDisconnectResponse.getReasonString()); |
|||
} |
|||
|
|||
@Override |
|||
public void mqttErrorOccurred(MqttException e) { |
|||
log.warn("Error occurred:", e); |
|||
} |
|||
|
|||
@Override |
|||
public void messageArrived(String requestTopic, MqttMessage mqttMessage) { |
|||
lastReceivedMessage = mqttMessage; |
|||
if (awaitSubTopic == null) { |
|||
log.warn("messageArrived on topic: {}", requestTopic); |
|||
qoS = mqttMessage.getQos(); |
|||
payloadBytes = mqttMessage.getPayload(); |
|||
subscribeLatch.countDown(); |
|||
} else { |
|||
messageArrivedOnAwaitSubTopic(requestTopic, mqttMessage); |
|||
} |
|||
} |
|||
|
|||
protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) { |
|||
log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); |
|||
if (awaitSubTopic.equals(requestTopic)) { |
|||
qoS = mqttMessage.getQos(); |
|||
payloadBytes = mqttMessage.getPayload(); |
|||
subscribeLatch.countDown(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void deliveryComplete(IMqttToken iMqttToken) { |
|||
log.warn("delivery complete: {}", iMqttToken.getResponse()); |
|||
pubAckReceived = iMqttToken.getResponse().getType() == MqttWireMessage.MESSAGE_TYPE_PUBACK; |
|||
deliveryLatch.countDown(); |
|||
} |
|||
|
|||
@Override |
|||
public void connectComplete(boolean reconnect, String serverURI) { |
|||
log.warn("Connect completed: reconnect - {}, serverURI - {}", reconnect, serverURI); |
|||
} |
|||
|
|||
@Override |
|||
public void authPacketArrived(int reasonCode, MqttProperties mqttProperties) { |
|||
log.warn("Auth package received: reasonCode - {}, mqtt properties - {}", reasonCode, mqttProperties); |
|||
} |
|||
} |
|||
@ -0,0 +1,175 @@ |
|||
/** |
|||
* 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.transport.mqtt.mqttv5; |
|||
|
|||
import io.netty.handler.codec.mqtt.MqttQoS; |
|||
import org.eclipse.paho.mqttv5.client.IMqttToken; |
|||
import org.eclipse.paho.mqttv5.client.MqttAsyncClient; |
|||
import org.eclipse.paho.mqttv5.client.MqttCallback; |
|||
import org.eclipse.paho.mqttv5.client.MqttConnectionOptions; |
|||
import org.eclipse.paho.mqttv5.client.persist.MemoryPersistence; |
|||
import org.eclipse.paho.mqttv5.common.MqttException; |
|||
import org.eclipse.paho.mqttv5.common.MqttMessage; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
|
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
public class MqttV5TestClient { // We should copy part of MqttV3TestClient, due to different package names in import
|
|||
|
|||
private static final String MQTT_URL = "tcp://localhost:1883"; |
|||
private static final int TIMEOUT = 30; // seconds
|
|||
private static final long TIMEOUT_MS = TimeUnit.SECONDS.toMillis(TIMEOUT); |
|||
|
|||
private final MqttAsyncClient client; |
|||
|
|||
public void setCallback(MqttCallback callback) { |
|||
client.setCallback(callback); |
|||
} |
|||
|
|||
public MqttV5TestClient() throws MqttException { |
|||
this.client = createClient(); |
|||
} |
|||
|
|||
public MqttV5TestClient(String clientId) throws MqttException { |
|||
this.client = createClient(clientId); |
|||
} |
|||
|
|||
public MqttV5TestClient(boolean generateClientId) throws MqttException { |
|||
this.client = createClient(generateClientId); |
|||
} |
|||
|
|||
public IMqttToken connectAndWait(String userName, String password) throws MqttException { |
|||
IMqttToken connect = connect(userName, password); |
|||
connect.waitForCompletion(TIMEOUT_MS); |
|||
return connect; |
|||
} |
|||
|
|||
public IMqttToken connectAndWait(String userName) throws MqttException { |
|||
return connectAndWait(userName, null); |
|||
} |
|||
|
|||
public IMqttToken connectAndWait() throws MqttException { |
|||
return connectAndWait(null, null); |
|||
} |
|||
|
|||
public IMqttToken connectAndWait(MqttConnectionOptions options) throws MqttException { |
|||
IMqttToken iMqttToken = connect(options); |
|||
iMqttToken.waitForCompletion(TIMEOUT_MS); |
|||
return iMqttToken; |
|||
} |
|||
|
|||
private IMqttToken connect(String userName, String password) throws MqttException { |
|||
if (client == null) { |
|||
throw new RuntimeException("Failed to connect! MqttAsyncClient is not initialized!"); |
|||
} |
|||
MqttConnectionOptions options = new MqttConnectionOptions(); |
|||
if (StringUtils.isNotEmpty(userName)) { |
|||
options.setUserName(userName); |
|||
} |
|||
if (StringUtils.isNotEmpty(password)) { |
|||
options.setPassword(password.getBytes()); |
|||
} |
|||
return client.connect(options); |
|||
} |
|||
|
|||
public IMqttToken connect(MqttConnectionOptions options) throws MqttException { |
|||
if (client == null) { |
|||
throw new RuntimeException("Failed to connect! MqttAsyncClient is not initialized!"); |
|||
} |
|||
return client.connect(options); |
|||
} |
|||
|
|||
public void disconnectAndWait() throws MqttException { |
|||
disconnect().waitForCompletion(TIMEOUT_MS); |
|||
} |
|||
|
|||
public IMqttToken disconnect() throws MqttException { |
|||
return client.disconnect(); |
|||
} |
|||
|
|||
public void disconnectForcibly() throws MqttException { |
|||
client.disconnectForcibly(TIMEOUT_MS); |
|||
} |
|||
|
|||
public IMqttToken publishAndWait(String topic, byte[] payload) throws MqttException { |
|||
IMqttToken iMqttToken = publish(topic, payload); |
|||
iMqttToken.waitForCompletion(TIMEOUT_MS); |
|||
return iMqttToken; |
|||
} |
|||
|
|||
public IMqttToken publish(String topic, byte[] payload) throws MqttException { |
|||
MqttMessage message = new MqttMessage(); |
|||
message.setPayload(payload); |
|||
return publish(topic, message); |
|||
} |
|||
|
|||
public IMqttToken publish(String topic, MqttMessage message) throws MqttException { |
|||
return publish(topic, message.getPayload(), message.getQos(), message.isRetained()); |
|||
} |
|||
|
|||
public IMqttToken publish(String topic, byte[] payload, int qos, boolean retain) throws MqttException { |
|||
return client.publish(topic, payload, qos, retain); |
|||
} |
|||
|
|||
public IMqttToken subscribeAndWait(String topic, MqttQoS qoS) throws MqttException { |
|||
IMqttToken iMqttToken = subscribe(topic, qoS); |
|||
iMqttToken.waitForCompletion(TIMEOUT_MS); |
|||
return iMqttToken; |
|||
} |
|||
|
|||
public IMqttToken subscribe(String topic, MqttQoS qoS) throws MqttException { |
|||
return client.subscribe(topic, qoS.value()); |
|||
} |
|||
|
|||
public IMqttToken unsubscribeAndWait(String topic) throws MqttException { |
|||
IMqttToken iMqttToken = unsubscribe(topic); |
|||
iMqttToken.waitForCompletion(TIMEOUT_MS); |
|||
return iMqttToken; |
|||
} |
|||
|
|||
public IMqttToken unsubscribe(String topic) throws MqttException { |
|||
return client.unsubscribe(topic); |
|||
} |
|||
|
|||
public boolean isConnected() { |
|||
return client.isConnected(); |
|||
} |
|||
|
|||
public void enableManualAcks() { |
|||
client.setManualAcks(true); |
|||
} |
|||
|
|||
public void messageArrivedComplete(MqttMessage mqttMessage) throws MqttException { |
|||
client.messageArrivedComplete(mqttMessage.getId(), mqttMessage.getQos()); |
|||
} |
|||
|
|||
private MqttAsyncClient createClient() throws MqttException { |
|||
return createClient(true); |
|||
} |
|||
|
|||
private MqttAsyncClient createClient(boolean generateClientId) throws MqttException { |
|||
String clientId = null; |
|||
if (generateClientId) { |
|||
clientId = "test" + System.nanoTime(); |
|||
} |
|||
return createClient(clientId); |
|||
} |
|||
|
|||
private MqttAsyncClient createClient(String clientId) throws MqttException { |
|||
return new MqttAsyncClient(MQTT_URL, clientId, new MemoryPersistence()); |
|||
} |
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue