From 89fa943b097897726ec8573b0631c54b06c0d9a4 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 20 May 2026 17:39:10 +0300 Subject: [PATCH] Add Kafka OAUTHBEARER (OAuth2 client-credentials) support for js-executor and tb-core --- .../src/main/resources/thingsboard.yml | 10 +++ .../server/queue/kafka/TbKafkaSettings.java | 24 ++++++- .../TbKafkaSettingsConfluentPlainTest.java | 52 ++++++++++++++ .../queue/kafka/TbKafkaSettingsOAuthTest.java | 59 +++++++++++++++ msa/js-executor/package.json | 2 +- msa/js-executor/queue/oAuthBearerProvider.ts | 72 ++++++++++++------- 6 files changed, 190 insertions(+), 29 deletions(-) create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6ceef7bc2e..ce1c931961 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1782,11 +1782,21 @@ queue: # The endpoint identification algorithm used by clients to validate server hostname. The default value is https ssl.algorithm: "${TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM:https}" # The mechanism used to authenticate Schema Registry requests. SASL/PLAIN should only be used with TLS/SSL as a transport layer to ensure that clear passwords are not transmitted on the wire without encryption + # Set to OAUTHBEARER to authenticate via OAuth2/OIDC client-credentials (configure the oauth.* settings below) sasl.mechanism: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM:PLAIN}" # Using JAAS Configuration for specifying multiple SASL mechanisms on a broker sasl.config: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG:org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";}" # Protocol used to communicate with brokers. Valid values are: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL security.protocol: "${TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL:SASL_SSL}" + # Optional: OAuth2 (OIDC) client-credentials settings. Applied only when sasl.mechanism is OAUTHBEARER. + # ThingsBoard wires Kafka's built-in OAuthBearerLoginCallbackHandler; the token is fetched and refreshed automatically. + oauth: + # Optional: OAuth2 client id + client-id: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID:}" + # Optional: OAuth2 client secret + client-secret: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET:}" + # Optional: OAuth2/OIDC token endpoint URL to request the bearer token from + endpoint-url: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL:}" # Key-value properties for Kafka consumer per specific topic, e.g. tb_ota_package is a topic name for ota, tb_rule_engine.sq is a topic name for default SequentialByOriginator queue. # Check TB_QUEUE_CORE_OTA_TOPIC and TB_QUEUE_RE_SQ_TOPIC params consumer-properties-per-topic: 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 fd9f469eff..797e975a54 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 @@ -23,6 +23,7 @@ import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.config.SaslConfigs; import org.apache.kafka.common.config.SslConfigs; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -137,6 +138,15 @@ public class TbKafkaSettings { @Value("${queue.kafka.confluent.security.protocol:}") private String securityProtocol; + @Value("${queue.kafka.confluent.oauth.client-id:}") + private String oauthClientId; + + @Value("${queue.kafka.confluent.oauth.client-secret:}") + private String oauthClientSecret; + + @Value("${queue.kafka.confluent.oauth.endpoint-url:}") + private String oauthEndpointUrl; + @Value("${queue.kafka.other-inline:}") private String otherInline; @@ -213,9 +223,19 @@ public class TbKafkaSettings { if (useConfluent) { props.put("ssl.endpoint.identification.algorithm", sslAlgorithm); - props.put("sasl.mechanism", saslMechanism); - props.put("sasl.jaas.config", saslConfig); + 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); + } else { + props.put(SaslConfigs.SASL_JAAS_CONFIG, saslConfig); + } } props.put(CommonClientConfigs.REQUEST_TIMEOUT_MS_CONFIG, requestTimeoutMs); diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java new file mode 100644 index 0000000000..4414756610 --- /dev/null +++ b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java @@ -0,0 +1,52 @@ +/** + * Copyright © 2016-2026 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.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=PLAIN", + "queue.kafka.confluent.security.protocol=SASL_SSL", + "queue.kafka.confluent.sasl.config=org.apache.kafka.common.security.plain.PlainLoginModule required username=\"key\" password=\"secret\";" +}) +class TbKafkaSettingsConfluentPlainTest { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenPlainMechanism_whenToProps_thenUsesLegacySaslConfigAndNoOauthKeys() { + Properties props = settings.toProps(); + + assertThat(props).containsEntry("sasl.mechanism", "PLAIN"); + assertThat(props.getProperty("sasl.jaas.config")) + .isEqualTo("org.apache.kafka.common.security.plain.PlainLoginModule required username=\"key\" password=\"secret\";"); + assertThat(props).doesNotContainKey("sasl.login.callback.handler.class"); + assertThat(props).doesNotContainKey("sasl.oauthbearer.token.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/TbKafkaSettingsOAuthTest.java new file mode 100644 index 0000000000..242290395f --- /dev/null +++ b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java @@ -0,0 +1,59 @@ +/** + * Copyright © 2016-2026 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.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=my-client", + "queue.kafka.confluent.oauth.client-secret=my-secret", + "queue.kafka.confluent.oauth.endpoint-url=https://idp.example.com/oauth/token" +}) +class TbKafkaSettingsOAuthTest { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenOauthbearerMechanism_whenToProps_thenConfiguresNativeOidcHandler() { + Properties props = settings.toProps(); + + assertThat(props).containsEntry("sasl.mechanism", "OAUTHBEARER"); + assertThat(props).containsEntry("security.protocol", "SASL_SSL"); + assertThat(props).containsEntry("sasl.login.callback.handler.class", + "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler"); + assertThat(props).containsEntry("sasl.oauthbearer.token.endpoint.url", + "https://idp.example.com/oauth/token"); + assertThat(props.getProperty("sasl.jaas.config")) + .contains("org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required") + .contains("clientId=\"my-client\"") + .contains("clientSecret=\"my-secret\""); + } + +} diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 8dc9b171b4..508d214219 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -14,7 +14,6 @@ }, "dependencies": { "@2l/kafkajs-lz4": "^1.3.2", - "@types/simple-oauth2": "^5.0.8", "config": "^4.1.1", "express": "^5.1.0", "js-yaml": "^4.1.1", @@ -37,6 +36,7 @@ "@types/config": "^3.3.5", "@types/express": "~5.0.3", "@types/node": "~22.17.2", + "@types/simple-oauth2": "^5.0.8", "@types/uuid-parse": "^1.0.2", "@yao-pkg/pkg": "^6.6.0", "fs-extra": "^11.3.1", diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index 08323d2046..7e56a5779d 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -14,9 +14,8 @@ /// limitations under the License. /// -import { AccessToken, ClientCredentials } from 'simple-oauth2' +import { AccessToken, ClientCredentials } from 'simple-oauth2'; import { _logger } from '../config/logger'; -import { error } from 'winston'; interface OauthBearerProviderOptions { clientId: string; @@ -25,8 +24,11 @@ interface OauthBearerProviderOptions { refreshThresholdMs: number; } +const RETRY_DELAY_MS = 5000; +const MIN_REFRESH_DELAY_MS = 1000; + export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { - const logger = _logger('oauthBearerProvider') + const logger = _logger('oauthBearerProvider'); const client = new ClientCredentials({ client: { id: options.clientId, @@ -38,38 +40,56 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { }); let tokenPromise: Promise; - let accessToken: AccessToken; + let refreshTimer: NodeJS.Timeout | undefined; - async function refreshToken():Promise{ - logger.info('Start token refreshing/validation'); + function scheduleRefresh(delayMs: number): void { + if (refreshTimer) { + clearTimeout(refreshTimer); + } + const delay = Math.max(delayMs, MIN_REFRESH_DELAY_MS); + logger.info('Next Kafka OAuth token refresh in %dms', delay); + refreshTimer = setTimeout(() => { + tokenPromise = refreshToken(); + // refreshToken() reschedules its own retry on failure; swallow here to avoid an unhandled rejection. + tokenPromise.catch((err) => { + logger.error('Scheduled Kafka OAuth token refresh failed: %s', err?.message ?? err); + }); + }, delay); + refreshTimer.unref(); + } + + async function refreshToken(): Promise { + logger.info('Requesting Kafka OAuth bearer token'); try { - if (accessToken == null) { - accessToken = await client.getToken({}) - logger.info('Got new token'); - } + // Client-credentials grant issues no refresh_token, so always request a fresh token. + const accessToken: AccessToken = await client.getToken({}); - if (accessToken.expired(options.refreshThresholdMs / 1000)) { - logger.info(`Token will expire during next ${options.refreshThresholdMs}ms. Refresh token`); - accessToken = await accessToken.refresh() + const expiresIn = accessToken.token.expires_in; + if (typeof expiresIn !== 'number' || expiresIn <= 0) { + throw new Error(`OAuth token response has invalid "expires_in": ${expiresIn}`); } - let expires_in = typeof accessToken.token.expires_in === "number" ? accessToken.token.expires_in : 0; //throw exception - const nextRefresh = expires_in * 1000 - options.refreshThresholdMs; - logger.info(`Next token validation in ${nextRefresh}ms.`); - setTimeout(() => { - tokenPromise = refreshToken() - }, nextRefresh); - let access_token = typeof accessToken.token.access_token === "string" ? accessToken.token.access_token : ""; //throw exception - return access_token; - } catch (error) { - throw error; + const accessTokenValue = accessToken.token.access_token; + if (typeof accessTokenValue !== 'string' || accessTokenValue.length === 0) { + throw new Error('OAuth token response has no "access_token"'); + } + + const nextRefresh = expiresIn * 1000 - options.refreshThresholdMs; + scheduleRefresh(nextRefresh); + return accessTokenValue; + } catch (err) { + logger.error('Failed to obtain Kafka OAuth bearer token, retrying in %dms', RETRY_DELAY_MS); + scheduleRefresh(RETRY_DELAY_MS); + throw err; } } tokenPromise = refreshToken(); + // The first fetch is retried internally; swallow here, so a startup failure is not an unhandled rejection. + tokenPromise.catch(() => { /* retry already scheduled in refreshToken() */ }); return async function () { return { value: await tokenPromise - } - } -}; \ No newline at end of file + }; + }; +};