Browse Source
# Conflicts: # application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java # application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.javapull/7511/head
367 changed files with 12711 additions and 5556 deletions
@ -0,0 +1,37 @@ |
|||
/** |
|||
* 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.install; |
|||
|
|||
import org.springframework.context.annotation.Profile; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.dao.util.NoSqlAnyDaoNonCloud; |
|||
|
|||
/* |
|||
* Create keyspace for Cassandra NoSQL database for non-cloud deployment. |
|||
* For cloud service like Astra DBaas admin have to create keyspace manually on cloud UI. |
|||
* Then create tokens with database admin role and put it on Thingsboard parameters. |
|||
* Without this service cloud DB will end up with exception like |
|||
* UnauthorizedException: Missing correct permission on thingsboard |
|||
* */ |
|||
@Service |
|||
@NoSqlAnyDaoNonCloud |
|||
@Profile("install") |
|||
public class CassandraKeyspaceService extends CassandraAbstractDatabaseSchemaService |
|||
implements NoSqlKeyspaceService { |
|||
public CassandraKeyspaceService() { |
|||
super("schema-keyspace.cql"); |
|||
} |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* 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.install; |
|||
|
|||
import org.springframework.context.annotation.Profile; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.service.executors.DbCallbackExecutorService; |
|||
|
|||
@Component |
|||
@Profile("install") |
|||
public class DbUpgradeExecutorService extends DbCallbackExecutorService { |
|||
|
|||
} |
|||
@ -0,0 +1,19 @@ |
|||
/** |
|||
* 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.install; |
|||
|
|||
public interface NoSqlKeyspaceService extends DatabaseSchemaService { |
|||
} |
|||
@ -0,0 +1,71 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
import io.jsonwebtoken.Claims; |
|||
import lombok.RequiredArgsConstructor; |
|||
import org.springframework.beans.factory.annotation.Qualifier; |
|||
import org.springframework.context.event.EventListener; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.cache.TbTransactionalCache; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.id.UserId; |
|||
import org.thingsboard.server.common.data.security.event.UserAuthDataChangedEvent; |
|||
import org.thingsboard.server.common.data.security.model.JwtToken; |
|||
import org.thingsboard.server.service.security.model.token.JwtTokenFactory; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
import static java.util.concurrent.TimeUnit.MILLISECONDS; |
|||
|
|||
@Service |
|||
public class DefaultTokenOutdatingService implements TokenOutdatingService { |
|||
|
|||
private final TbTransactionalCache<String, Long> cache; |
|||
private final JwtTokenFactory tokenFactory; |
|||
|
|||
public DefaultTokenOutdatingService(@Qualifier("UsersSessionInvalidation") TbTransactionalCache<String, Long> cache, JwtTokenFactory tokenFactory) { |
|||
this.cache = cache; |
|||
this.tokenFactory = tokenFactory; |
|||
} |
|||
|
|||
@EventListener(classes = UserAuthDataChangedEvent.class) |
|||
public void onUserAuthDataChanged(UserAuthDataChangedEvent event) { |
|||
if (StringUtils.hasText(event.getId())) { |
|||
cache.put(event.getId(), event.getTs()); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public boolean isOutdated(JwtToken token, UserId userId) { |
|||
Claims claims = tokenFactory.parseTokenClaims(token).getBody(); |
|||
long issueTime = claims.getIssuedAt().getTime(); |
|||
String sessionId = claims.get("sessionId", String.class); |
|||
if (isTokenOutdated(issueTime, userId.toString())){ |
|||
return true; |
|||
} else { |
|||
return sessionId != null && isTokenOutdated(issueTime, sessionId); |
|||
} |
|||
} |
|||
|
|||
private Boolean isTokenOutdated(long issueTime, String sessionId) { |
|||
return Optional.ofNullable(cache.get(sessionId)).map(outdatageTime -> isTokenOutdated(issueTime, outdatageTime.get())).orElse(false); |
|||
} |
|||
|
|||
private boolean isTokenOutdated(long issueTime, Long outdatageTime) { |
|||
return MILLISECONDS.toSeconds(issueTime) < MILLISECONDS.toSeconds(outdatageTime); |
|||
} |
|||
} |
|||
@ -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,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); |
|||
} |
|||
} |
|||
@ -0,0 +1,34 @@ |
|||
/** |
|||
* 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.cache.usersUpdateTime; |
|||
|
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
|||
import org.springframework.cache.CacheManager; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.cache.CaffeineTbTransactionalCache; |
|||
import org.thingsboard.server.common.data.CacheConstants; |
|||
|
|||
|
|||
@ConditionalOnProperty(prefix = "cache", value = "type", havingValue = "caffeine", matchIfMissing = true) |
|||
@Service("UsersSessionInvalidation") |
|||
public class UsersSessionInvalidationCaffeineCache extends CaffeineTbTransactionalCache<String, Long> { |
|||
|
|||
@Autowired |
|||
public UsersSessionInvalidationCaffeineCache(CacheManager cacheManager) { |
|||
super(cacheManager, CacheConstants.USERS_SESSION_INVALIDATION_CACHE); |
|||
} |
|||
} |
|||
@ -0,0 +1,36 @@ |
|||
/** |
|||
* 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.cache.usersUpdateTime; |
|||
|
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
|||
import org.springframework.data.redis.connection.RedisConnectionFactory; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.cache.CacheSpecsMap; |
|||
import org.thingsboard.server.cache.RedisTbTransactionalCache; |
|||
import org.thingsboard.server.cache.TBRedisCacheConfiguration; |
|||
import org.thingsboard.server.cache.TbFSTRedisSerializer; |
|||
import org.thingsboard.server.common.data.CacheConstants; |
|||
|
|||
@ConditionalOnProperty(prefix = "cache", value = "type", havingValue = "redis") |
|||
@Service("UsersSessionInvalidation") |
|||
public class UsersSessionInvalidationRedisCache extends RedisTbTransactionalCache<String, Long> { |
|||
|
|||
@Autowired |
|||
public UsersSessionInvalidationRedisCache(TBRedisCacheConfiguration configuration, CacheSpecsMap cacheSpecsMap, RedisConnectionFactory connectionFactory) { |
|||
super(CacheConstants.USERS_SESSION_INVALIDATION_CACHE, cacheSpecsMap, connectionFactory, configuration, new TbFSTRedisSerializer<>()); |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue