Browse Source

Refactoring after review

pull/15670/head
Andrii Landiak 5 months ago
parent
commit
5af6bef320
  1. 31
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  2. 71
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java
  3. 63
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java
  4. 2
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java
  5. 30
      msa/js-executor/queue/kafkaTemplate.ts
  6. 12
      msa/js-executor/queue/oAuthBearerProvider.ts

31
common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java

@ -33,6 +33,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TbProperty;
import org.thingsboard.server.queue.util.PropertyUtils;
import org.thingsboard.server.queue.util.TbKafkaComponent;
@ -226,13 +227,7 @@ public class TbKafkaSettings {
props.put(SaslConfigs.SASL_MECHANISM, saslMechanism);
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
if ("OAUTHBEARER".equalsIgnoreCase(saslMechanism)) {
props.put(SaslConfigs.SASL_JAAS_CONFIG,
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required"
+ " clientId=\"" + oauthClientId + "\""
+ " clientSecret=\"" + oauthClientSecret + "\";");
props.put(SaslConfigs.SASL_LOGIN_CALLBACK_HANDLER_CLASS,
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler");
props.put(SaslConfigs.SASL_OAUTHBEARER_TOKEN_ENDPOINT_URL, oauthEndpointUrl);
applyOauthBearerProps(props);
} else {
props.put(SaslConfigs.SASL_JAAS_CONFIG, saslConfig);
}
@ -250,6 +245,28 @@ public class TbKafkaSettings {
return props;
}
private void applyOauthBearerProps(Properties props) {
if (StringUtils.isBlank(oauthClientId) || StringUtils.isBlank(oauthClientSecret) || StringUtils.isBlank(oauthEndpointUrl)) {
throw new IllegalStateException("Kafka SASL mechanism is OAUTHBEARER but "
+ "queue.kafka.confluent.oauth.client-id / client-secret / endpoint-url are not all set");
}
if (!oauthEndpointUrl.toLowerCase().startsWith("https://")) {
log.warn("Kafka OAuth token endpoint URL is not HTTPS ({}); client credentials will be sent unencrypted",
oauthEndpointUrl);
}
props.put(SaslConfigs.SASL_JAAS_CONFIG,
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required"
+ " clientId=\"" + escapeJaasValue(oauthClientId) + "\""
+ " clientSecret=\"" + escapeJaasValue(oauthClientSecret) + "\";");
props.put(SaslConfigs.SASL_LOGIN_CALLBACK_HANDLER_CLASS,
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler");
props.put(SaslConfigs.SASL_OAUTHBEARER_TOKEN_ENDPOINT_URL, oauthEndpointUrl);
}
private static String escapeJaasValue(String value) {
return value.replace("\\", "\\\\").replace("\"", "\\\"");
}
void configureSSL(Properties props) {
if (sslEnabled) {
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL");

71
common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java

@ -0,0 +1,71 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2026 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.queue.kafka;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.TestPropertySource;
import java.util.Properties;
import static org.assertj.core.api.Assertions.assertThat;
@SpringBootTest(classes = {TbKafkaSettings.class, KafkaAdmin.class})
@TestPropertySource(properties = {
"queue.type=kafka",
"queue.kafka.bootstrap.servers=localhost:9092",
"queue.kafka.use_confluent_cloud=true",
"queue.kafka.confluent.sasl.mechanism=OAUTHBEARER",
"queue.kafka.confluent.security.protocol=SASL_SSL",
"queue.kafka.confluent.oauth.client-id=client\"with\\\\quote",
"queue.kafka.confluent.oauth.client-secret=secret\"with\\\\backslash;and-semi",
"queue.kafka.confluent.oauth.endpoint-url=https://idp.example.com/oauth/token"
})
class TbKafkaSettingsConfluentOAuthEscapingTest {
@Autowired
TbKafkaSettings settings;
// Locks in that double quotes and backslashes in client credentials are escaped
// so the JAAS config string parses correctly and cannot inject extra options.
@Test
void givenSecretWithQuoteAndBackslash_whenToProps_thenEscapesJaasValues() {
Properties props = settings.toProps();
String jaas = props.getProperty("sasl.jaas.config");
assertThat(jaas)
.contains("clientId=\"client\\\"with\\\\quote\"")
.contains("clientSecret=\"secret\\\"with\\\\backslash;and-semi\"")
.endsWith("\";");
}
}

63
common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java

@ -0,0 +1,63 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2026 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.queue.kafka;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.TestPropertySource;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@SpringBootTest(classes = {TbKafkaSettings.class, KafkaAdmin.class})
@TestPropertySource(properties = {
"queue.type=kafka",
"queue.kafka.bootstrap.servers=localhost:9092",
"queue.kafka.use_confluent_cloud=true",
"queue.kafka.confluent.sasl.mechanism=OAUTHBEARER",
"queue.kafka.confluent.security.protocol=SASL_SSL",
"queue.kafka.confluent.oauth.client-id=my-client",
"queue.kafka.confluent.oauth.client-secret=my-secret"
})
class TbKafkaSettingsConfluentOAuthMissingConfigTest {
@Autowired
TbKafkaSettings settings;
@Test
void givenOauthbearerWithoutEndpointUrl_whenToProps_thenFailsFast() {
assertThatThrownBy(() -> settings.toProps())
.isInstanceOf(IllegalStateException.class)
.hasMessageContaining("OAUTHBEARER")
.hasMessageContaining("endpoint-url");
}
}

2
common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java → common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java

@ -35,7 +35,7 @@ import static org.assertj.core.api.Assertions.assertThat;
"queue.kafka.confluent.oauth.client-secret=my-secret",
"queue.kafka.confluent.oauth.endpoint-url=https://idp.example.com/oauth/token"
})
class TbKafkaSettingsOAuthTest {
class TbKafkaSettingsConfluentOAuthTest {
@Autowired
TbKafkaSettings settings;

30
msa/js-executor/queue/kafkaTemplate.ts

@ -19,7 +19,7 @@ import fs from 'node:fs';
import { _logger, KafkaJsWinstonLogCreator } from '../config/logger';
import { JsInvokeMessageProcessor } from '../api/jsInvokeMessageProcessor'
import { IQueue } from './queue.models';
import {oauthBearerProvider} from './oAuthBearerProvider';
import { oauthBearerProvider } from './oAuthBearerProvider';
import {
Admin,
CompressionCodecs,
@ -37,6 +37,8 @@ import { KeyObject } from 'tls';
import process, { exit, kill } from 'process';
const DEFAULT_OAUTH_REFRESH_THRESHOLD_MS = 60000;
export class KafkaTemplate implements IQueue {
private logger = _logger(`kafkaTemplate`);
@ -115,21 +117,23 @@ export class KafkaTemplate implements IQueue {
kafkaConfig['connectionTimeout'] = this.connectionTimeout;
if (useConfluent) {
let saslMechanism = config.get('kafka.confluent.sasl.mechanism') as String
if(saslMechanism.toLowerCase() == "oauthbearer"){
const saslMechanism = config.get('kafka.confluent.sasl.mechanism') as string;
if (saslMechanism.toLowerCase() === 'oauthbearer') {
const refreshThresholdMs = config.has('kafka.confluent.oauth.refresh_threshold')
? Number(config.get('kafka.confluent.oauth.refresh_threshold'))
: DEFAULT_OAUTH_REFRESH_THRESHOLD_MS;
kafkaConfig['sasl'] = {
mechanism: "oauthbearer",
oauthBearerProvider: oauthBearerProvider({
clientId: config.get('kafka.confluent.oauth.client_id'),
clientSecret: config.get('kafka.confluent.oauth.client_secret'),
host: config.get('kafka.confluent.oauth.endpoint_url'),
refreshThresholdMs: config.has('kafka.confluent.oauth.refresh_threshold')?config.get('kafka.confluent.oauth.refresh_threshold'):60000,
})
mechanism: 'oauthbearer',
oauthBearerProvider: oauthBearerProvider({
clientId: config.get('kafka.confluent.oauth.client_id'),
clientSecret: config.get('kafka.confluent.oauth.client_secret'),
host: config.get('kafka.confluent.oauth.endpoint_url'),
refreshThresholdMs,
})
};
}
else{
} else {
kafkaConfig['sasl'] = {
mechanism: config.get('kafka.confluent.sasl.mechanism') as any,
mechanism: saslMechanism as any,
username: config.get('kafka.confluent.username'),
password: config.get('kafka.confluent.password')
};

12
msa/js-executor/queue/oAuthBearerProvider.ts

@ -29,6 +29,16 @@ const MIN_REFRESH_DELAY_MS = 1000;
export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
const logger = _logger('oauthBearerProvider');
if (!options.clientId || !options.clientSecret || !options.host) {
throw new Error('Kafka OAUTHBEARER requires kafka.confluent.oauth.client_id, client_secret and endpoint_url to be set');
}
if (!/^https:\/\//i.test(options.host)) {
logger.warn('Kafka OAuth token endpoint URL is not HTTPS (%s); client credentials will be sent unencrypted', options.host);
}
const refreshThresholdMs = Number(options.refreshThresholdMs);
if (!Number.isFinite(refreshThresholdMs) || refreshThresholdMs < 0) {
throw new Error(`Kafka OAuth refresh_threshold must be a non-negative number, got: ${options.refreshThresholdMs}`);
}
const client = new ClientCredentials({
client: {
id: options.clientId,
@ -73,7 +83,7 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
throw new Error('OAuth token response has no "access_token"');
}
const nextRefresh = expiresIn * 1000 - options.refreshThresholdMs;
const nextRefresh = expiresIn * 1000 - refreshThresholdMs;
scheduleRefresh(nextRefresh);
return accessTokenValue;
} catch (err) {

Loading…
Cancel
Save