From ffbff7d88ed8349b7506a88daf5e48dcf2159164 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 4 Jun 2026 14:31:58 +0300 Subject: [PATCH] Add optional OAuth scope to Kafka OAUTHBEARER on tb-core path --- .../src/main/resources/thingsboard.yml | 2 + .../server/queue/kafka/TbKafkaSettings.java | 13 ++++- ...bKafkaSettingsConfluentOAuthScopeTest.java | 56 +++++++++++++++++++ .../config/custom-environment-variables.yml | 2 + msa/js-executor/queue/kafkaTemplate.ts | 4 ++ msa/js-executor/queue/oAuthBearerProvider.ts | 52 +++++++++++++---- 6 files changed, 115 insertions(+), 14 deletions(-) create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index ce1c931961..ca450417dd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1797,6 +1797,8 @@ queue: 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:}" + # Optional: OAuth2 scope/audience requested for the token. Required by some IdPs for client-credentials (e.g. Azure AD's "api:///.default") + scope: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_SCOPE:}" # 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 f2b70ff20e..d71657fcbb 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 @@ -148,6 +148,9 @@ public class TbKafkaSettings { @Value("${queue.kafka.confluent.oauth.endpoint-url:}") private String oauthEndpointUrl; + @Value("${queue.kafka.confluent.oauth.scope:}") + private String oauthScope; + @Value("${queue.kafka.other-inline:}") private String otherInline; @@ -254,10 +257,16 @@ public class TbKafkaSettings { log.warn("Kafka OAuth token endpoint URL is not HTTPS ({}); client credentials will be sent unencrypted", oauthEndpointUrl); } - props.put(SaslConfigs.SASL_JAAS_CONFIG, + StringBuilder jaasConfig = new StringBuilder( "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required" + " clientId=\"" + escapeJaasValue(oauthClientId) + "\"" - + " clientSecret=\"" + escapeJaasValue(oauthClientSecret) + "\";"); + + " clientSecret=\"" + escapeJaasValue(oauthClientSecret) + "\""); + if (StringUtils.isNotBlank(oauthScope)) { + // Some IdPs (e.g. Azure AD's ".default") require a scope for the client-credentials grant. + jaasConfig.append(" scope=\"").append(escapeJaasValue(oauthScope)).append("\""); + } + jaasConfig.append(";"); + props.put(SaslConfigs.SASL_JAAS_CONFIG, jaasConfig.toString()); props.put(SaslConfigs.SASL_LOGIN_CALLBACK_HANDLER_CLASS, "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler"); props.put(SaslConfigs.SASL_OAUTHBEARER_TOKEN_ENDPOINT_URL, oauthEndpointUrl); diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java new file mode 100644 index 0000000000..aa81d3131e --- /dev/null +++ b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java @@ -0,0 +1,56 @@ +/** + * 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", + "queue.kafka.confluent.oauth.scope=api://my-app/.default" +}) +class TbKafkaSettingsConfluentOAuthScopeTest { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenOauthScope_whenToProps_thenAppendsScopeToJaasConfig() { + Properties props = settings.toProps(); + + assertThat(props.getProperty("sasl.jaas.config")) + .contains("org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required") + .contains("clientId=\"my-client\"") + .contains("clientSecret=\"my-secret\"") + .contains("scope=\"api://my-app/.default\"") + .endsWith("\";"); + } + +} diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 5297ee3ad4..74c519b6fb 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -64,6 +64,8 @@ kafka: endpoint_url: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL" # Optional: threshold to renew current token before it expires (ms) refresh_threshold: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_REFRESH_THRESHOLD" + # Optional: Oauth scope/audience requested for the token. Required by some IdPs for client-credentials (e.g. Azure AD's "api:///.default") + scope: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_SCOPE" logger: level: "LOGGER_LEVEL" diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index e5c2ccd67f..d7b6e7896f 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -122,6 +122,9 @@ export class KafkaTemplate implements IQueue { const refreshThresholdMs = config.has('kafka.confluent.oauth.refresh_threshold') ? Number(config.get('kafka.confluent.oauth.refresh_threshold')) : DEFAULT_OAUTH_REFRESH_THRESHOLD_MS; + const scope = config.has('kafka.confluent.oauth.scope') + ? config.get('kafka.confluent.oauth.scope') as string + : undefined; kafkaConfig['sasl'] = { mechanism: 'oauthbearer', oauthBearerProvider: oauthBearerProvider({ @@ -129,6 +132,7 @@ export class KafkaTemplate implements IQueue { clientSecret: config.get('kafka.confluent.oauth.client_secret'), host: config.get('kafka.confluent.oauth.endpoint_url'), refreshThresholdMs, + scope, }) }; } else { diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index ba60e55482..f4bf5fe1b2 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -22,6 +22,8 @@ interface OauthBearerProviderOptions { clientSecret: string; host: string; refreshThresholdMs: number; + // Optional OAuth2 scope. Required by some IdPs for client-credentials (e.g. Azure AD's "api:///.default"). + scope?: string; } const RETRY_DELAY_MS = 5000; @@ -39,6 +41,7 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { if (!Number.isFinite(refreshThresholdMs) || refreshThresholdMs < 0) { throw new Error(`Kafka OAuth refresh_threshold must be a non-negative number, got: ${options.refreshThresholdMs}`); } + const scope = options.scope && options.scope.trim().length > 0 ? options.scope.trim() : undefined; let tokenUrl: URL; try { tokenUrl = new URL(options.host); @@ -59,8 +62,14 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { } }); - let tokenPromise: Promise; + // Last successfully fetched token. Kept across refreshes so a failed background refresh + // does not drop a still-valid token out from under KafkaJS (which calls this provider on + // each reauthentication). It is only swapped in after a successful refresh. + let cachedToken: string | undefined; + // Resolves once the very first token has been obtained (used only until cachedToken is set). + let initialToken: Promise; let refreshTimer: NodeJS.Timeout | undefined; + let warnedShortLivedToken = false; function scheduleRefresh(delayMs: number): void { if (refreshTimer) { @@ -69,9 +78,9 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { 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) => { + // refreshToken() updates cachedToken on success and reschedules its own retry on failure; + // swallow here so a failed background refresh is not an unhandled rejection. + refreshToken().catch((err) => { logger.error('Scheduled Kafka OAuth token refresh failed: %s', err?.message ?? err); }); }, delay); @@ -82,18 +91,36 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { logger.info('Requesting Kafka OAuth bearer token'); try { // Client-credentials grant issues no refresh_token, so always request a fresh token. - const accessToken: AccessToken = await client.getToken({}); + const accessToken: AccessToken = await client.getToken(scope ? { scope } : {}); - const expiresIn = accessToken.token.expires_in; - if (typeof expiresIn !== 'number' || expiresIn <= 0) { - throw new Error(`OAuth token response has invalid "expires_in": ${expiresIn}`); + // Some IdPs return expires_in as a numeric string ("3600"); coerce before validating. + const rawExpiresIn = accessToken.token.expires_in; + const expiresIn = Number(rawExpiresIn); + if (!Number.isFinite(expiresIn) || expiresIn <= 0) { + throw new Error(`OAuth token response has invalid "expires_in": ${rawExpiresIn}`); } 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 - refreshThresholdMs; + cachedToken = accessTokenValue; + + const lifetimeMs = expiresIn * 1000; + // If the token lifetime is shorter than the refresh threshold, refreshing "threshold + // before expiry" would fire immediately and clamp to MIN_REFRESH_DELAY_MS, hammering the + // token endpoint every second. Cap the effective threshold to half the lifetime instead. + let effectiveThresholdMs = refreshThresholdMs; + if (refreshThresholdMs >= lifetimeMs) { + effectiveThresholdMs = lifetimeMs / 2; + if (!warnedShortLivedToken) { + logger.warn('Kafka OAuth token lifetime (%dms) is <= refresh threshold (%dms); capping threshold to %dms to avoid hammering the token endpoint', + lifetimeMs, refreshThresholdMs, effectiveThresholdMs); + warnedShortLivedToken = true; + } + } + + const nextRefresh = lifetimeMs - effectiveThresholdMs; scheduleRefresh(nextRefresh); return accessTokenValue; } catch (err) { @@ -103,13 +130,14 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { } } - tokenPromise = refreshToken(); + initialToken = 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() */ }); + initialToken.catch(() => { /* retry already scheduled in refreshToken() */ }); return async function () { + // Prefer the last good token; only fall back to awaiting the initial fetch before one exists. return { - value: await tokenPromise + value: cachedToken ?? await initialToken }; }; };