Browse Source

Add Kafka OAUTHBEARER (OAuth2 client-credentials) support for js-executor and tb-core

pull/15670/head
Andrii Landiak 4 months ago
parent
commit
89fa943b09
  1. 10
      application/src/main/resources/thingsboard.yml
  2. 24
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  3. 52
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java
  4. 59
      common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsOAuthTest.java
  5. 2
      msa/js-executor/package.json
  6. 72
      msa/js-executor/queue/oAuthBearerProvider.ts

10
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:

24
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);

52
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");
}
}

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

2
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",

72
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<string>;
let accessToken: AccessToken;
let refreshTimer: NodeJS.Timeout | undefined;
async function refreshToken():Promise<string>{
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<string> {
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
}
}
};
};
};
};

Loading…
Cancel
Save