Browse Source
# Conflicts: # rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNodeConfiguration.javapull/7895/head
392 changed files with 20407 additions and 5260 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,162 @@ |
|||
/** |
|||
* 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.jwt.settings; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.RandomStringUtils; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.context.annotation.Lazy; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.cluster.TbClusterService; |
|||
import org.thingsboard.server.common.data.AdminSettings; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|||
import org.thingsboard.server.common.data.security.model.JwtSettings; |
|||
import org.thingsboard.server.dao.settings.AdminSettingsService; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
import java.util.Base64; |
|||
import java.util.Objects; |
|||
import java.util.Optional; |
|||
|
|||
@Service |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class DefaultJwtSettingsService implements JwtSettingsService { |
|||
|
|||
@Lazy |
|||
private final AdminSettingsService adminSettingsService; |
|||
@Lazy |
|||
private final Optional<TbClusterService> tbClusterService; |
|||
private final JwtSettingsValidator jwtSettingsValidator; |
|||
|
|||
@Value("${security.jwt.tokenExpirationTime:9000}") |
|||
private Integer tokenExpirationTime; |
|||
@Value("${security.jwt.refreshTokenExpTime:604800}") |
|||
private Integer refreshTokenExpTime; |
|||
@Value("${security.jwt.tokenIssuer:thingsboard.io}") |
|||
private String tokenIssuer; |
|||
@Value("${security.jwt.tokenSigningKey:thingsboardDefaultSigningKey}") |
|||
private String tokenSigningKey; |
|||
|
|||
private volatile JwtSettings jwtSettings = null; //lazy init
|
|||
|
|||
/** |
|||
* Create JWT admin settings is intended to be called from Install scripts only |
|||
*/ |
|||
@Override |
|||
public void createRandomJwtSettings() { |
|||
if (getJwtSettingsFromDb() == null) { |
|||
log.info("Creating JWT admin settings..."); |
|||
this.jwtSettings = getJwtSettingsFromYml(); |
|||
if (isSigningKeyDefault(jwtSettings)) { |
|||
this.jwtSettings.setTokenSigningKey(Base64.getEncoder().encodeToString( |
|||
RandomStringUtils.randomAlphanumeric(64).getBytes(StandardCharsets.UTF_8))); |
|||
} |
|||
saveJwtSettings(jwtSettings); |
|||
} else { |
|||
log.info("Skip creating JWT admin settings because they already exist."); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* Create JWT admin settings is intended to be called from Upgrade scripts only |
|||
*/ |
|||
@Override |
|||
public void saveLegacyYmlSettings() { |
|||
log.info("Saving legacy JWT admin settings from YML..."); |
|||
if (getJwtSettingsFromDb() == null) { |
|||
saveJwtSettings(getJwtSettingsFromYml()); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public JwtSettings saveJwtSettings(JwtSettings jwtSettings) { |
|||
jwtSettingsValidator.validate(jwtSettings); |
|||
final AdminSettings adminJwtSettings = mapJwtToAdminSettings(jwtSettings); |
|||
final AdminSettings existedSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, ADMIN_SETTINGS_JWT_KEY); |
|||
if (existedSettings != null) { |
|||
adminJwtSettings.setId(existedSettings.getId()); |
|||
} |
|||
|
|||
log.info("Saving new JWT admin settings. From this moment, the JWT parameters from YAML and ENV will be ignored"); |
|||
adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, adminJwtSettings); |
|||
|
|||
tbClusterService.ifPresent(cs -> cs.broadcastEntityStateChangeEvent(TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID, ComponentLifecycleEvent.UPDATED)); |
|||
return reloadJwtSettings(); |
|||
} |
|||
|
|||
@Override |
|||
public JwtSettings reloadJwtSettings() { |
|||
return getJwtSettings(true); |
|||
} |
|||
|
|||
@Override |
|||
public JwtSettings getJwtSettings() { |
|||
return getJwtSettings(false); |
|||
} |
|||
|
|||
public JwtSettings getJwtSettings(boolean forceReload) { |
|||
if (this.jwtSettings == null || forceReload) { |
|||
synchronized (this) { |
|||
if (this.jwtSettings == null || forceReload) { |
|||
JwtSettings result = getJwtSettingsFromDb(); |
|||
if (result == null) { |
|||
result = getJwtSettingsFromYml(); |
|||
log.warn("Loading the JWT settings from YML since there are no settings in DB. Looks like the upgrade script was not applied."); |
|||
} |
|||
if (isSigningKeyDefault(result)) { |
|||
log.warn("WARNING: The platform is configured to use default JWT Signing Key. " + |
|||
"This is a security issue that needs to be resolved. Please change the JWT Signing Key using the Web UI. " + |
|||
"Navigate to \"System settings -> Security settings\" while logged in as a System Administrator."); |
|||
} |
|||
this.jwtSettings = result; |
|||
} |
|||
} |
|||
} |
|||
return this.jwtSettings; |
|||
} |
|||
|
|||
private JwtSettings getJwtSettingsFromYml() { |
|||
return new JwtSettings(this.tokenExpirationTime, this.refreshTokenExpTime, this.tokenIssuer, this.tokenSigningKey); |
|||
} |
|||
|
|||
private JwtSettings getJwtSettingsFromDb() { |
|||
AdminSettings adminJwtSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, ADMIN_SETTINGS_JWT_KEY); |
|||
return adminJwtSettings != null ? mapAdminToJwtSettings(adminJwtSettings) : null; |
|||
} |
|||
|
|||
private JwtSettings mapAdminToJwtSettings(AdminSettings adminSettings) { |
|||
Objects.requireNonNull(adminSettings, "adminSettings for JWT is null"); |
|||
return JacksonUtil.treeToValue(adminSettings.getJsonValue(), JwtSettings.class); |
|||
} |
|||
|
|||
private AdminSettings mapJwtToAdminSettings(JwtSettings jwtSettings) { |
|||
Objects.requireNonNull(jwtSettings, "jwtSettings is null"); |
|||
AdminSettings adminJwtSettings = new AdminSettings(); |
|||
adminJwtSettings.setTenantId(TenantId.SYS_TENANT_ID); |
|||
adminJwtSettings.setKey(ADMIN_SETTINGS_JWT_KEY); |
|||
adminJwtSettings.setJsonValue(JacksonUtil.valueToTree(jwtSettings)); |
|||
return adminJwtSettings; |
|||
} |
|||
|
|||
private boolean isSigningKeyDefault(JwtSettings settings) { |
|||
return TOKEN_SIGNING_KEY_DEFAULT.equals(settings.getTokenSigningKey()); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,69 @@ |
|||
/** |
|||
* 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.jwt.settings; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.apache.commons.lang3.RandomUtils; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.bouncycastle.util.Arrays; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.security.model.JwtSettings; |
|||
import org.thingsboard.server.dao.exception.DataValidationException; |
|||
|
|||
import java.util.Base64; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@Component |
|||
@RequiredArgsConstructor |
|||
public class DefaultJwtSettingsValidator implements JwtSettingsValidator { |
|||
|
|||
@Override |
|||
public void validate(JwtSettings jwtSettings) { |
|||
if (StringUtils.isEmpty(jwtSettings.getTokenIssuer())) { |
|||
throw new DataValidationException("JWT token issuer should be specified!"); |
|||
} |
|||
if (Optional.ofNullable(jwtSettings.getRefreshTokenExpTime()).orElse(0) <= TimeUnit.MINUTES.toSeconds(15)) { |
|||
throw new DataValidationException("JWT refresh token expiration time should be at least 15 minutes!"); |
|||
} |
|||
if (Optional.ofNullable(jwtSettings.getTokenExpirationTime()).orElse(0) <= TimeUnit.MINUTES.toSeconds(1)) { |
|||
throw new DataValidationException("JWT token expiration time should be at least 1 minute!"); |
|||
} |
|||
if (jwtSettings.getTokenExpirationTime() >= jwtSettings.getRefreshTokenExpTime()) { |
|||
throw new DataValidationException("JWT token expiration time should greater than JWT refresh token expiration time!"); |
|||
} |
|||
if (StringUtils.isEmpty(jwtSettings.getTokenSigningKey())) { |
|||
throw new DataValidationException("JWT token signing key should be specified!"); |
|||
} |
|||
|
|||
byte[] decodedKey; |
|||
try { |
|||
decodedKey = Base64.getDecoder().decode(jwtSettings.getTokenSigningKey()); |
|||
} catch (Exception e) { |
|||
throw new DataValidationException("JWT token signing key should be a valid Base64 encoded string! " + e.getMessage()); |
|||
} |
|||
|
|||
if (Arrays.isNullOrEmpty(decodedKey)) { |
|||
throw new DataValidationException("JWT token signing key should be non-empty after Base64 decoding!"); |
|||
} |
|||
if (decodedKey.length * Byte.SIZE < 256 && !JwtSettingsService.TOKEN_SIGNING_KEY_DEFAULT.equals(jwtSettings.getTokenSigningKey())) { |
|||
throw new DataValidationException("JWT token signing key should be a Base64 encoded string representing at least 256 bits of data!"); |
|||
} |
|||
|
|||
System.arraycopy(decodedKey, 0, RandomUtils.nextBytes(decodedKey.length), 0, decodedKey.length); //secure memory
|
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* 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.jwt.settings; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.springframework.context.annotation.Primary; |
|||
import org.springframework.context.annotation.Profile; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.security.model.JwtSettings; |
|||
|
|||
/** |
|||
* During Install or upgrade the validation is suppressed to keep existing data |
|||
* */ |
|||
|
|||
@Primary |
|||
@Profile("install") |
|||
@Component |
|||
@RequiredArgsConstructor |
|||
public class InstallJwtSettingsValidator implements JwtSettingsValidator { |
|||
|
|||
@Override |
|||
public void validate(JwtSettings jwtSettings) { |
|||
|
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,35 @@ |
|||
/** |
|||
* 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.jwt.settings; |
|||
|
|||
import org.thingsboard.server.common.data.security.model.JwtSettings; |
|||
|
|||
public interface JwtSettingsService { |
|||
|
|||
String ADMIN_SETTINGS_JWT_KEY = "jwt"; |
|||
String TOKEN_SIGNING_KEY_DEFAULT = "thingsboardDefaultSigningKey"; |
|||
|
|||
JwtSettings getJwtSettings(); |
|||
|
|||
JwtSettings reloadJwtSettings(); |
|||
|
|||
void createRandomJwtSettings(); |
|||
|
|||
void saveLegacyYmlSettings(); |
|||
|
|||
JwtSettings saveJwtSettings(JwtSettings jwtSettings); |
|||
|
|||
} |
|||
@ -0,0 +1,23 @@ |
|||
/** |
|||
* 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.jwt.settings; |
|||
|
|||
import org.thingsboard.server.common.data.security.model.JwtSettings; |
|||
|
|||
public interface JwtSettingsValidator { |
|||
|
|||
void validate(JwtSettings jwtSettings); |
|||
} |
|||
@ -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); |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue