From 5af6bef32079d7a8805c0275ea069574c4c85963 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 21 May 2026 11:24:58 +0300 Subject: [PATCH] Refactoring after review --- .../server/queue/kafka/TbKafkaSettings.java | 31 ++++++-- ...fkaSettingsConfluentOAuthEscapingTest.java | 71 +++++++++++++++++++ ...ttingsConfluentOAuthMissingConfigTest.java | 63 ++++++++++++++++ ...=> TbKafkaSettingsConfluentOAuthTest.java} | 2 +- msa/js-executor/queue/kafkaTemplate.ts | 30 ++++---- msa/js-executor/queue/oAuthBearerProvider.ts | 12 +++- 6 files changed, 187 insertions(+), 22 deletions(-) create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java rename common/queue/src/test/java/org/thingsboard/server/queue/kafka/{TbKafkaSettingsOAuthTest.java => TbKafkaSettingsConfluentOAuthTest.java} (98%) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java index 797e975a54..a676395420 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java +++ b/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"); diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java new file mode 100644 index 0000000000..ef3e8fed19 --- /dev/null +++ b/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("\";"); + } + +} diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java new file mode 100644 index 0000000000..092ef4ca2f --- /dev/null +++ b/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"); + } + +} diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java similarity index 98% rename from common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java rename to common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java index 242290395f..4e56fd60fc 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java +++ b/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; diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index ab9896dd8e..e5c2ccd67f 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/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') }; diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index 7e56a5779d..9d45986390 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/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) {