Browse Source

Add optional OAuth scope to Kafka OAUTHBEARER on tb-core path

pull/15670/head
Andrii Landiak 4 months ago
parent
commit
ffbff7d88e
  1. 2
      application/src/main/resources/thingsboard.yml
  2. 13
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  3. 56
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java
  4. 2
      msa/js-executor/config/custom-environment-variables.yml
  5. 4
      msa/js-executor/queue/kafkaTemplate.ts
  6. 52
      msa/js-executor/queue/oAuthBearerProvider.ts

2
application/src/main/resources/thingsboard.yml

@ -1797,6 +1797,8 @@ queue:
client-secret: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET:}" client-secret: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET:}"
# Optional: OAuth2/OIDC token endpoint URL to request the bearer token from # Optional: OAuth2/OIDC token endpoint URL to request the bearer token from
endpoint-url: "${TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL:}" 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://<id>/.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. # 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 # Check TB_QUEUE_CORE_OTA_TOPIC and TB_QUEUE_RE_SQ_TOPIC params
consumer-properties-per-topic: consumer-properties-per-topic:

13
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:}") @Value("${queue.kafka.confluent.oauth.endpoint-url:}")
private String oauthEndpointUrl; private String oauthEndpointUrl;
@Value("${queue.kafka.confluent.oauth.scope:}")
private String oauthScope;
@Value("${queue.kafka.other-inline:}") @Value("${queue.kafka.other-inline:}")
private String otherInline; 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", log.warn("Kafka OAuth token endpoint URL is not HTTPS ({}); client credentials will be sent unencrypted",
oauthEndpointUrl); oauthEndpointUrl);
} }
props.put(SaslConfigs.SASL_JAAS_CONFIG, StringBuilder jaasConfig = new StringBuilder(
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required" "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required"
+ " clientId=\"" + escapeJaasValue(oauthClientId) + "\"" + " 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, props.put(SaslConfigs.SASL_LOGIN_CALLBACK_HANDLER_CLASS,
"org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler"); "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler");
props.put(SaslConfigs.SASL_OAUTHBEARER_TOKEN_ENDPOINT_URL, oauthEndpointUrl); props.put(SaslConfigs.SASL_OAUTHBEARER_TOKEN_ENDPOINT_URL, oauthEndpointUrl);

56
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("\";");
}
}

2
msa/js-executor/config/custom-environment-variables.yml

@ -64,6 +64,8 @@ kafka:
endpoint_url: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL" endpoint_url: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL"
# Optional: threshold to renew current token before it expires (ms) # Optional: threshold to renew current token before it expires (ms)
refresh_threshold: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_REFRESH_THRESHOLD" 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://<id>/.default")
scope: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_SCOPE"
logger: logger:
level: "LOGGER_LEVEL" level: "LOGGER_LEVEL"

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

@ -122,6 +122,9 @@ export class KafkaTemplate implements IQueue {
const refreshThresholdMs = config.has('kafka.confluent.oauth.refresh_threshold') const refreshThresholdMs = config.has('kafka.confluent.oauth.refresh_threshold')
? Number(config.get('kafka.confluent.oauth.refresh_threshold')) ? Number(config.get('kafka.confluent.oauth.refresh_threshold'))
: DEFAULT_OAUTH_REFRESH_THRESHOLD_MS; : DEFAULT_OAUTH_REFRESH_THRESHOLD_MS;
const scope = config.has('kafka.confluent.oauth.scope')
? config.get('kafka.confluent.oauth.scope') as string
: undefined;
kafkaConfig['sasl'] = { kafkaConfig['sasl'] = {
mechanism: 'oauthbearer', mechanism: 'oauthbearer',
oauthBearerProvider: oauthBearerProvider({ oauthBearerProvider: oauthBearerProvider({
@ -129,6 +132,7 @@ export class KafkaTemplate implements IQueue {
clientSecret: config.get('kafka.confluent.oauth.client_secret'), clientSecret: config.get('kafka.confluent.oauth.client_secret'),
host: config.get('kafka.confluent.oauth.endpoint_url'), host: config.get('kafka.confluent.oauth.endpoint_url'),
refreshThresholdMs, refreshThresholdMs,
scope,
}) })
}; };
} else { } else {

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

@ -22,6 +22,8 @@ interface OauthBearerProviderOptions {
clientSecret: string; clientSecret: string;
host: string; host: string;
refreshThresholdMs: number; refreshThresholdMs: number;
// Optional OAuth2 scope. Required by some IdPs for client-credentials (e.g. Azure AD's "api://<id>/.default").
scope?: string;
} }
const RETRY_DELAY_MS = 5000; const RETRY_DELAY_MS = 5000;
@ -39,6 +41,7 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
if (!Number.isFinite(refreshThresholdMs) || refreshThresholdMs < 0) { if (!Number.isFinite(refreshThresholdMs) || refreshThresholdMs < 0) {
throw new Error(`Kafka OAuth refresh_threshold must be a non-negative number, got: ${options.refreshThresholdMs}`); 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; let tokenUrl: URL;
try { try {
tokenUrl = new URL(options.host); tokenUrl = new URL(options.host);
@ -59,8 +62,14 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
} }
}); });
let tokenPromise: Promise<string>; // 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<string>;
let refreshTimer: NodeJS.Timeout | undefined; let refreshTimer: NodeJS.Timeout | undefined;
let warnedShortLivedToken = false;
function scheduleRefresh(delayMs: number): void { function scheduleRefresh(delayMs: number): void {
if (refreshTimer) { if (refreshTimer) {
@ -69,9 +78,9 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
const delay = Math.max(delayMs, MIN_REFRESH_DELAY_MS); const delay = Math.max(delayMs, MIN_REFRESH_DELAY_MS);
logger.info('Next Kafka OAuth token refresh in %dms', delay); logger.info('Next Kafka OAuth token refresh in %dms', delay);
refreshTimer = setTimeout(() => { refreshTimer = setTimeout(() => {
tokenPromise = refreshToken(); // refreshToken() updates cachedToken on success and reschedules its own retry on failure;
// refreshToken() reschedules its own retry on failure; swallow here to avoid an unhandled rejection. // swallow here so a failed background refresh is not an unhandled rejection.
tokenPromise.catch((err) => { refreshToken().catch((err) => {
logger.error('Scheduled Kafka OAuth token refresh failed: %s', err?.message ?? err); logger.error('Scheduled Kafka OAuth token refresh failed: %s', err?.message ?? err);
}); });
}, delay); }, delay);
@ -82,18 +91,36 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => {
logger.info('Requesting Kafka OAuth bearer token'); logger.info('Requesting Kafka OAuth bearer token');
try { try {
// Client-credentials grant issues no refresh_token, so always request a fresh token. // 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; // Some IdPs return expires_in as a numeric string ("3600"); coerce before validating.
if (typeof expiresIn !== 'number' || expiresIn <= 0) { const rawExpiresIn = accessToken.token.expires_in;
throw new Error(`OAuth token response has invalid "expires_in": ${expiresIn}`); 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; const accessTokenValue = accessToken.token.access_token;
if (typeof accessTokenValue !== 'string' || accessTokenValue.length === 0) { if (typeof accessTokenValue !== 'string' || accessTokenValue.length === 0) {
throw new Error('OAuth token response has no "access_token"'); 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); scheduleRefresh(nextRefresh);
return accessTokenValue; return accessTokenValue;
} catch (err) { } 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. // 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 () { return async function () {
// Prefer the last good token; only fall back to awaiting the initial fetch before one exists.
return { return {
value: await tokenPromise value: cachedToken ?? await initialToken
}; };
}; };
}; };

Loading…
Cancel
Save