From f7e356cd3620e23bbf3cdfd245d475569c4d1367 Mon Sep 17 00:00:00 2001 From: Jonas Koch Date: Tue, 24 Feb 2026 15:41:07 +0100 Subject: [PATCH 01/15] js-executor kafka add oauthbearer provider --- .../config/custom-environment-variables.yml | 5 ++ msa/js-executor/package.json | 2 + msa/js-executor/queue/kafkaTemplate.ts | 25 ++++-- msa/js-executor/queue/oAuthBearerProvider.ts | 67 ++++++++++++++ msa/js-executor/yarn.lock | 88 +++++++++++++++++++ 5 files changed, 182 insertions(+), 5 deletions(-) create mode 100644 msa/js-executor/queue/oAuthBearerProvider.ts diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 6a98361c72..4cf12fc1f6 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -54,6 +54,11 @@ kafka: mechanism: "TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM" username: "TB_QUEUE_KAFKA_CONFLUENT_USERNAME" password: "TB_QUEUE_KAFKA_CONFLUENT_PASSWORD" + oauth: + client_id: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID" + client_secret: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET" + endpoint_url: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL" + refresh_threshold: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_REFRESH_THRESHOLD" logger: level: "LOGGER_LEVEL" diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index ed4565ec08..8a27d13007 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -13,11 +13,13 @@ "build": "tsc" }, "dependencies": { + "@types/simple-oauth2": "^5.0.8", "config": "^4.1.1", "express": "^5.1.0", "js-yaml": "^4.1.1", "kafkajs": "^2.2.4", "long": "^5.3.2", + "simple-oauth2": "^5.1.0", "uuid-parse": "^1.1.0", "winston": "^3.17.0", "winston-daily-rotate-file": "^5.0.0" diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index 4527e9e1ce..3be36a8061 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -19,6 +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 { Admin, CompressionTypes, @@ -89,11 +90,25 @@ export class KafkaTemplate implements IQueue { kafkaConfig['connectionTimeout'] = this.connectionTimeout; if (useConfluent) { - kafkaConfig['sasl'] = { - mechanism: config.get('kafka.confluent.sasl.mechanism') as any, - username: config.get('kafka.confluent.username'), - password: config.get('kafka.confluent.password') - }; + let saslMechanism = config.get('kafka.confluent.sasl.mechanism') as String + if(saslMechanism.toLowerCase() == "oauthbearer"){ + 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.get('kafka.confluent.oauth.refresh_threshold'), + }) + }; + } + else{ + kafkaConfig['sasl'] = { + mechanism: config.get('kafka.confluent.sasl.mechanism') as any, + username: config.get('kafka.confluent.username'), + password: config.get('kafka.confluent.password') + }; + } kafkaConfig['ssl'] = true; } diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts new file mode 100644 index 0000000000..2312fbc054 --- /dev/null +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -0,0 +1,67 @@ +/// +/// 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. +/// + +import { AccessToken, ClientCredentials } from 'simple-oauth2' +interface OauthBearerProviderOptions { + clientId: string; + clientSecret: string; + host: string; + refreshThresholdMs: number; +} + +export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { + const client = new ClientCredentials({ + client: { + id: options.clientId, + secret: options.clientSecret + }, + auth: { + tokenHost: options.host + } + }); + + let tokenPromise: Promise; + let accessToken: AccessToken; + + async function refreshToken():Promise{ + try { + if (accessToken == null) { + accessToken = await client.getToken({}) + } + + if (accessToken.expired(options.refreshThresholdMs / 1000)) { + accessToken = await accessToken.refresh() + } + let expires_in = typeof accessToken.token.expires_in === "number" ? accessToken.token.expires_in : 0; //throw exception + const nextRefresh = expires_in * 1000 - options.refreshThresholdMs; + 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; + } + } + + tokenPromise = refreshToken(); + + return async function () { + return { + value: await tokenPromise + } + } +}; \ No newline at end of file diff --git a/msa/js-executor/yarn.lock b/msa/js-executor/yarn.lock index 1e4a802d3c..535c88964a 100644 --- a/msa/js-executor/yarn.lock +++ b/msa/js-executor/yarn.lock @@ -59,6 +59,44 @@ enabled "2.0.x" kuler "^2.0.0" +"@hapi/boom@^10.0.1": + version "10.0.1" + resolved "https://registry.yarnpkg.com/@hapi/boom/-/boom-10.0.1.tgz#ebb14688275ae150aa6af788dbe482e6a6062685" + integrity sha512-ERcCZaEjdH3OgSJlyjVk8pHIFeus91CjKP3v+MpgBNp5IvGzP2l/bRiD78nqYcKPaZdbKkK5vDBVPd2ohHBlsA== + dependencies: + "@hapi/hoek" "^11.0.2" + +"@hapi/bourne@^3.0.0": + version "3.0.0" + resolved "https://registry.yarnpkg.com/@hapi/bourne/-/bourne-3.0.0.tgz#f11fdf7dda62fe8e336fa7c6642d9041f30356d7" + integrity sha512-Waj1cwPXJDucOib4a3bAISsKJVb15MKi9IvmTI/7ssVEm6sywXGjVJDhl6/umt1pK1ZS7PacXU3A1PmFKHEZ2w== + +"@hapi/hoek@^11.0.2", "@hapi/hoek@^11.0.4": + version "11.0.7" + resolved "https://registry.yarnpkg.com/@hapi/hoek/-/hoek-11.0.7.tgz#56a920793e0a42d10e530da9a64cc0d3919c4002" + integrity sha512-HV5undWkKzcB4RZUusqOpcgxOaq6VOAH7zhhIr2g3G8NF/MlFO75SjOr2NfuSx0Mh40+1FqCkagKLJRykUWoFQ== + +"@hapi/hoek@^9.0.0", "@hapi/hoek@^9.3.0": + version "9.3.0" + resolved "https://registry.yarnpkg.com/@hapi/hoek/-/hoek-9.3.0.tgz#8368869dcb735be2e7f5cb7647de78e167a251fb" + integrity sha512-/c6rf4UJlmHlC9b5BaNvzAcFv7HZ2QHaV0D4/HNlBdvFnvQq8RI4kYdhyPCl7Xj+oWvTWQ8ujhqS53LIgAe6KQ== + +"@hapi/topo@^5.1.0": + version "5.1.0" + resolved "https://registry.yarnpkg.com/@hapi/topo/-/topo-5.1.0.tgz#dc448e332c6c6e37a4dc02fd84ba8d44b9afb012" + integrity sha512-foQZKJig7Ob0BMAYBfcJk8d77QtOe7Wo4ox7ff1lQYoNNAb6jwcY1ncdoy2e9wQZzvNy7ODZCYJkK8kzmcAnAg== + dependencies: + "@hapi/hoek" "^9.0.0" + +"@hapi/wreck@^18.0.0": + version "18.1.0" + resolved "https://registry.yarnpkg.com/@hapi/wreck/-/wreck-18.1.0.tgz#68e631fc7568ebefc6252d5b86cb804466c8dbe6" + integrity sha512-0z6ZRCmFEfV/MQqkQomJ7sl/hyxvcZM7LtuVqN3vdAO4vM9eBbowl0kaqQj9EJJQab+3Uuh1GxbGIBFy4NfJ4w== + dependencies: + "@hapi/boom" "^10.0.1" + "@hapi/bourne" "^3.0.0" + "@hapi/hoek" "^11.0.2" + "@isaacs/fs-minipass@^4.0.0": version "4.0.1" resolved "https://registry.yarnpkg.com/@isaacs/fs-minipass/-/fs-minipass-4.0.1.tgz#2d59ae3ab4b38fb4270bfa23d30f8e2e86c7fe32" @@ -100,6 +138,23 @@ "@jridgewell/resolve-uri" "^3.1.0" "@jridgewell/sourcemap-codec" "^1.4.14" +"@sideway/address@^4.1.5": + version "4.1.5" + resolved "https://registry.yarnpkg.com/@sideway/address/-/address-4.1.5.tgz#4bc149a0076623ced99ca8208ba780d65a99b9d5" + integrity sha512-IqO/DUQHUkPeixNQ8n0JA6102hT9CmaljNTPmQ1u8MEhBo/R4Q8eKLN/vGZxuebwOroDB4cbpjheD4+/sKFK4Q== + dependencies: + "@hapi/hoek" "^9.0.0" + +"@sideway/formula@^3.0.1": + version "3.0.1" + resolved "https://registry.yarnpkg.com/@sideway/formula/-/formula-3.0.1.tgz#80fcbcbaf7ce031e0ef2dd29b1bfc7c3f583611f" + integrity sha512-/poHZJJVjx3L+zVD6g9KgHfYnb443oi7wLu/XKojDviHy6HOEOA6z1Trk5aR1dGcmPenJEgb2sK2I80LeS3MIg== + +"@sideway/pinpoint@^2.0.0": + version "2.0.0" + resolved "https://registry.yarnpkg.com/@sideway/pinpoint/-/pinpoint-2.0.0.tgz#cff8ffadc372ad29fd3f78277aeb29e632cc70df" + integrity sha512-RNiOoTPkptFtSVzQevY/yWtZwf/RxyVnPy/OcA9HBM3MlGDnBEYL5B41H0MTn0Uec8Hi+2qUtTfG2WWZBmMejQ== + "@tsconfig/node10@^1.0.7": version "1.0.11" resolved "https://registry.yarnpkg.com/@tsconfig/node10/-/node10-1.0.11.tgz#6ee46400685f130e278128c7b38b7e031ff5b2f2" @@ -210,6 +265,11 @@ "@types/node" "*" "@types/send" "*" +"@types/simple-oauth2@^5.0.8": + version "5.0.8" + resolved "https://registry.yarnpkg.com/@types/simple-oauth2/-/simple-oauth2-5.0.8.tgz#ec0c51a7df267e5d41b0b063aecdd34a312e9310" + integrity sha512-TehQqoOGdy3/rmFsCEGgnt1f4JhUCA0joWemGGCTbVYvoZvfBjkRsBFYmz8k0V/sn2XQZHe33L4lWxqMhIO3tQ== + "@types/triple-beam@^1.3.2": version "1.3.5" resolved "https://registry.yarnpkg.com/@types/triple-beam/-/triple-beam-1.3.5.tgz#74fef9ffbaa198eb8b588be029f38b00299caa2c" @@ -539,6 +599,13 @@ debug@4, debug@^4, debug@^4.3.5, debug@^4.4.0: dependencies: ms "^2.1.3" +debug@^4.3.4: + version "4.4.3" + resolved "https://registry.yarnpkg.com/debug/-/debug-4.4.3.tgz#c6ae432d9bd9662582fce08709b038c58e9e3d6a" + integrity sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA== + dependencies: + ms "^2.1.3" + decompress-response@^6.0.0: version "6.0.0" resolved "https://registry.yarnpkg.com/decompress-response/-/decompress-response-6.0.0.tgz#ca387612ddb7e104bd16d85aab00d5ecf09c66fc" @@ -945,6 +1012,17 @@ isarray@~1.0.0: resolved "https://registry.yarnpkg.com/isarray/-/isarray-1.0.0.tgz#bb935d48582cba168c06834957a54a3e07124f11" integrity sha512-VLghIWNM6ELQzo7zwmcg0NmTVyWKYjvIeM83yjp0wRDTmUnrM678fQbcKBo6n2CJEF0szoG//ytg+TKla89ALQ== +joi@^17.6.4: + version "17.13.3" + resolved "https://registry.yarnpkg.com/joi/-/joi-17.13.3.tgz#0f5cc1169c999b30d344366d384b12d92558bcec" + integrity sha512-otDA4ldcIx+ZXsKHWmp0YizCweVRZG96J10b0FevjfuncLO1oX59THoAmHkNubYJ+9gWsYsp5k8v4ib6oDv1fA== + dependencies: + "@hapi/hoek" "^9.3.0" + "@hapi/topo" "^5.1.0" + "@sideway/address" "^4.1.5" + "@sideway/formula" "^3.0.1" + "@sideway/pinpoint" "^2.0.0" + js-yaml@^4.1.1: version "4.1.1" resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-4.1.1.tgz#854c292467705b699476e1a2decc0c8a3458806b" @@ -1444,6 +1522,16 @@ simple-get@^4.0.0: once "^1.3.1" simple-concat "^1.0.0" +simple-oauth2@^5.1.0: + version "5.1.0" + resolved "https://registry.yarnpkg.com/simple-oauth2/-/simple-oauth2-5.1.0.tgz#1398fe2b8f4b4066298d63c155501b31b42238f2" + integrity sha512-gWDa38Ccm4MwlG5U7AlcJxPv3lvr80dU7ARJWrGdgvOKyzSj1gr3GBPN1rABTedAYvC/LsGYoFuFxwDBPtGEbw== + dependencies: + "@hapi/hoek" "^11.0.4" + "@hapi/wreck" "^18.0.0" + debug "^4.3.4" + joi "^17.6.4" + simple-swizzle@^0.2.2: version "0.2.2" resolved "https://registry.yarnpkg.com/simple-swizzle/-/simple-swizzle-0.2.2.tgz#a4da6b635ffcccca33f70d17cb92592de95e557a" From 896e0ba7b5b821be772173aefc3eaea6b20b3d6a Mon Sep 17 00:00:00 2001 From: Jonas Koch Date: Tue, 24 Feb 2026 17:03:41 +0100 Subject: [PATCH 02/15] adding description to config yaml --- msa/js-executor/config/custom-environment-variables.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 4cf12fc1f6..9e759f2b2f 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -54,10 +54,15 @@ kafka: mechanism: "TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM" username: "TB_QUEUE_KAFKA_CONFLUENT_USERNAME" password: "TB_QUEUE_KAFKA_CONFLUENT_PASSWORD" + # Optional: OAUTH Setting oauth: + # Optional: Oauth Client ID client_id: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID" + # Optional: Oauth Client Secret client_secret: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET" + # Optional: Oauth Endpoint URL to get the token from 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" logger: From 1d1c0437a0744a1f2323fa3ce6fdc83365af5379 Mon Sep 17 00:00:00 2001 From: Jonas Koch Date: Tue, 3 Mar 2026 15:42:42 +0100 Subject: [PATCH 03/15] add logging for token refresh --- msa/js-executor/queue/oAuthBearerProvider.ts | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index 2312fbc054..f09c5b2e51 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -15,6 +15,8 @@ /// import { AccessToken, ClientCredentials } from 'simple-oauth2' +import { _logger, KafkaJsWinstonLogCreator } from '../config/logger'; + interface OauthBearerProviderOptions { clientId: string; clientSecret: string; @@ -23,6 +25,7 @@ interface OauthBearerProviderOptions { } export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { + const logger = _logger('oauthBearerProvider') const client = new ClientCredentials({ client: { id: options.clientId, @@ -37,16 +40,20 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { let accessToken: AccessToken; async function refreshToken():Promise{ + logger.info('Start token refreshing/validation'); try { if (accessToken == null) { accessToken = await client.getToken({}) + logger.info('Got new token'); } if (accessToken.expired(options.refreshThresholdMs / 1000)) { + logger.info(`Token will expire during next ${options.refreshThresholdMs}ms. Refresh token`); accessToken = await accessToken.refresh() } 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); From 0b6cef9c5338891695c7cb1e067217e4b0f72de3 Mon Sep 17 00:00:00 2001 From: Jonas Koch Date: Tue, 19 May 2026 11:54:43 +0200 Subject: [PATCH 04/15] handle default values --- msa/js-executor/queue/kafkaTemplate.ts | 2 +- msa/js-executor/queue/oAuthBearerProvider.ts | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index 3be36a8061..30a53304d3 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -98,7 +98,7 @@ export class KafkaTemplate implements IQueue { 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.get('kafka.confluent.oauth.refresh_threshold'), + refreshThresholdMs: config.has('kafka.confluent.oauth.refresh_threshold')?config.get('kafka.confluent.oauth.refresh_threshold'):60000, }) }; } diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index f09c5b2e51..08323d2046 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -15,7 +15,8 @@ /// import { AccessToken, ClientCredentials } from 'simple-oauth2' -import { _logger, KafkaJsWinstonLogCreator } from '../config/logger'; +import { _logger } from '../config/logger'; +import { error } from 'winston'; interface OauthBearerProviderOptions { clientId: string; From 89fa943b097897726ec8573b0631c54b06c0d9a4 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 20 May 2026 17:39:10 +0300 Subject: [PATCH 05/15] 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 + }; + }; +}; From cca70df6e059fd6bd1225564d7746540fbfede3a Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Wed, 20 May 2026 18:12:37 +0300 Subject: [PATCH 06/15] yarn:install --- msa/js-executor/yarn.lock | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/msa/js-executor/yarn.lock b/msa/js-executor/yarn.lock index c045e6dbe2..c11a29f1b7 100644 --- a/msa/js-executor/yarn.lock +++ b/msa/js-executor/yarn.lock @@ -161,9 +161,9 @@ "@hapi/hoek" "^9.0.0" "@hapi/wreck@^18.0.0": - version "18.1.0" - resolved "https://registry.yarnpkg.com/@hapi/wreck/-/wreck-18.1.0.tgz#68e631fc7568ebefc6252d5b86cb804466c8dbe6" - integrity sha512-0z6ZRCmFEfV/MQqkQomJ7sl/hyxvcZM7LtuVqN3vdAO4vM9eBbowl0kaqQj9EJJQab+3Uuh1GxbGIBFy4NfJ4w== + version "18.1.2" + resolved "https://registry.yarnpkg.com/@hapi/wreck/-/wreck-18.1.2.tgz#1c1c84427085d7018ff4157ef3eb1e2cf080b26b" + integrity sha512-3dMnV2pfhQiyEqu8DL3VBmxkdLiRDiiUDuG79Dp+UK1gL9ZxAfDOUhB6k3D5MLqcgJJ1IARyGFhwoc1NITr/pg== dependencies: "@hapi/boom" "^10.0.1" "@hapi/bourne" "^3.0.0" From 5af6bef32079d7a8805c0275ea069574c4c85963 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 21 May 2026 11:24:58 +0300 Subject: [PATCH 07/15] 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) { From 89279b87c125e5e6a826446b9a7581858d58e297 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 21 May 2026 11:52:46 +0300 Subject: [PATCH 08/15] Refactoring --- .../server/queue/kafka/TbKafkaSettings.java | 2 +- ...fkaSettingsConfluentOAuthEscapingTest.java | 35 ++++++------------- ...ttingsConfluentOAuthMissingConfigTest.java | 35 ++++++------------- msa/js-executor/queue/oAuthBearerProvider.ts | 12 ++++++- 4 files changed, 32 insertions(+), 52 deletions(-) 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 a676395420..f2b70ff20e 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 @@ -250,7 +250,7 @@ public class TbKafkaSettings { 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://")) { + if (!oauthEndpointUrl.regionMatches(true, 0, "https://", 0, "https://".length())) { log.warn("Kafka OAuth token endpoint URL is not HTTPS ({}); client credentials will be sent unencrypted", oauthEndpointUrl); } 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 index ef3e8fed19..36857f8a2f 100644 --- 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 @@ -1,32 +1,17 @@ /** - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2026 The Thingsboard Authors * - * Copyright © 2016-2026 ThingsBoard, Inc. All Rights Reserved. + * 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 * - * 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. + * http://www.apache.org/licenses/LICENSE-2.0 * - * 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. + * 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; 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 index 092ef4ca2f..31f01f6121 100644 --- 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 @@ -1,32 +1,17 @@ /** - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2026 The Thingsboard Authors * - * Copyright © 2016-2026 ThingsBoard, Inc. All Rights Reserved. + * 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 * - * 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. + * http://www.apache.org/licenses/LICENSE-2.0 * - * 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. + * 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; diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index 9d45986390..ba60e55482 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -39,13 +39,23 @@ 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}`); } + let tokenUrl: URL; + try { + tokenUrl = new URL(options.host); + } catch { + throw new Error(`Kafka OAuth endpoint_url is not a valid URL: ${options.host}`); + } const client = new ClientCredentials({ client: { id: options.clientId, secret: options.clientSecret }, auth: { - tokenHost: options.host + // endpoint_url is the full token endpoint URL. Split it into host + path so + // simple-oauth2 does not append its default tokenPath (/oauth/token) and + // discard the real path (breaks Keycloak, Azure AD, Okta, etc.). + tokenHost: tokenUrl.origin, + tokenPath: tokenUrl.pathname + tokenUrl.search } }); From ffbff7d88ed8349b7506a88daf5e48dcf2159164 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 4 Jun 2026 14:31:58 +0300 Subject: [PATCH 09/15] 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 }; }; }; From d161ee0d9f6eb4589b25730d9f9d96dc2578bffb Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 4 Jun 2026 14:41:17 +0300 Subject: [PATCH 10/15] Changes after self-review --- msa/js-executor/queue/kafkaTemplate.ts | 10 +++-- msa/js-executor/queue/oAuthBearerProvider.ts | 46 ++++++++++++++++---- 2 files changed, 44 insertions(+), 12 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index d7b6e7896f..855b393e7b 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -125,12 +125,16 @@ export class KafkaTemplate implements IQueue { const scope = config.has('kafka.confluent.oauth.scope') ? config.get('kafka.confluent.oauth.scope') as string : undefined; + // Read as optional so a missing key yields '' (and the provider raises its clear + // "requires client_id, client_secret and endpoint_url" error) rather than node-config + // throwing a generic "Configuration property ... is not defined". + const optionalOauthStr = (key: string): string => config.has(key) ? config.get(key) as string : ''; 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'), + clientId: optionalOauthStr('kafka.confluent.oauth.client_id'), + clientSecret: optionalOauthStr('kafka.confluent.oauth.client_secret'), + host: optionalOauthStr('kafka.confluent.oauth.endpoint_url'), refreshThresholdMs, scope, }) diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index f4bf5fe1b2..899b6bd2e7 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -28,6 +28,10 @@ interface OauthBearerProviderOptions { const RETRY_DELAY_MS = 5000; const MIN_REFRESH_DELAY_MS = 1000; +// Safety margin applied at serve time: if the cached token is within this window of (or past) its +// hard expiry, fetch synchronously instead of serving it. Covers a scheduled refresh that slipped +// because the unref()'d timer was starved or the process was suspended. Also absorbs minor clock skew. +const EXPIRY_SAFETY_MS = 5000; export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { const logger = _logger('oauthBearerProvider'); @@ -66,8 +70,12 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { // 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; + // Epoch ms at which cachedToken hard-expires (access_token lifetime from the IdP). 0 until the + // first successful fetch. Used to decide at serve time whether the cached token is still safe. + let cachedTokenExpiresAt = 0; + // In-flight token fetch, shared so a serve-triggered refresh and the scheduled timer (and + // concurrent KafkaJS callbacks) coalesce onto one request instead of stampeding the IdP. + let refreshInFlight: Promise | undefined; let refreshTimer: NodeJS.Timeout | undefined; let warnedShortLivedToken = false; @@ -78,15 +86,26 @@ 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(() => { - // refreshToken() updates cachedToken on success and reschedules its own retry on failure; + // refreshTokenOnce() 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) => { + refreshTokenOnce().catch((err) => { logger.error('Scheduled Kafka OAuth token refresh failed: %s', err?.message ?? err); }); }, delay); refreshTimer.unref(); } + // Coalesce concurrent refreshes (scheduled timer, serve-time fallback, parallel KafkaJS callbacks) + // onto a single in-flight request so a near-expiry burst cannot stampede the token endpoint. + function refreshTokenOnce(): Promise { + if (!refreshInFlight) { + refreshInFlight = refreshToken().finally(() => { + refreshInFlight = undefined; + }); + } + return refreshInFlight; + } + async function refreshToken(): Promise { logger.info('Requesting Kafka OAuth bearer token'); try { @@ -104,9 +123,10 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { throw new Error('OAuth token response has no "access_token"'); } + const lifetimeMs = expiresIn * 1000; cachedToken = accessTokenValue; + cachedTokenExpiresAt = Date.now() + lifetimeMs; - 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. @@ -130,14 +150,22 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { } } - initialToken = refreshToken(); + // Kick off the initial fetch eagerly so the first authentication is fast. // The first fetch is retried internally; swallow here, so a startup failure is not an unhandled rejection. - initialToken.catch(() => { /* retry already scheduled in refreshToken() */ }); + refreshTokenOnce().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. + // Serve the cached token while it is comfortably valid; the background timer rotates it + // proactively. If there is no token yet, or the cached one is within EXPIRY_SAFETY_MS of its + // hard expiry (a scheduled refresh slipped), fetch synchronously so KafkaJS never reauthenticates + // with an expired token. Concurrent callers coalesce onto the same request via refreshTokenOnce(). + if (!cachedToken || Date.now() >= cachedTokenExpiresAt - EXPIRY_SAFETY_MS) { + return { + value: await refreshTokenOnce() + }; + } return { - value: cachedToken ?? await initialToken + value: cachedToken }; }; }; From a519e720f87f478c7246f51d16017ae5def6114a Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 4 Jun 2026 15:06:36 +0300 Subject: [PATCH 11/15] Renaming to endpointUrl for OauthBearerProviderOptions --- msa/js-executor/queue/kafkaTemplate.ts | 2 +- msa/js-executor/queue/oAuthBearerProvider.ts | 12 ++++++------ 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index 855b393e7b..95c95c6e39 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -134,7 +134,7 @@ export class KafkaTemplate implements IQueue { oauthBearerProvider: oauthBearerProvider({ clientId: optionalOauthStr('kafka.confluent.oauth.client_id'), clientSecret: optionalOauthStr('kafka.confluent.oauth.client_secret'), - host: optionalOauthStr('kafka.confluent.oauth.endpoint_url'), + endpointUrl: optionalOauthStr('kafka.confluent.oauth.endpoint_url'), refreshThresholdMs, scope, }) diff --git a/msa/js-executor/queue/oAuthBearerProvider.ts b/msa/js-executor/queue/oAuthBearerProvider.ts index 899b6bd2e7..83ae5125b6 100644 --- a/msa/js-executor/queue/oAuthBearerProvider.ts +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -20,7 +20,7 @@ import { _logger } from '../config/logger'; interface OauthBearerProviderOptions { clientId: string; clientSecret: string; - host: string; + endpointUrl: string; refreshThresholdMs: number; // Optional OAuth2 scope. Required by some IdPs for client-credentials (e.g. Azure AD's "api:///.default"). scope?: string; @@ -35,11 +35,11 @@ const EXPIRY_SAFETY_MS = 5000; export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { const logger = _logger('oauthBearerProvider'); - if (!options.clientId || !options.clientSecret || !options.host) { + if (!options.clientId || !options.clientSecret || !options.endpointUrl) { 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); + if (!/^https:\/\//i.test(options.endpointUrl)) { + logger.warn('Kafka OAuth token endpoint URL is not HTTPS (%s); client credentials will be sent unencrypted', options.endpointUrl); } const refreshThresholdMs = Number(options.refreshThresholdMs); if (!Number.isFinite(refreshThresholdMs) || refreshThresholdMs < 0) { @@ -48,9 +48,9 @@ export const oauthBearerProvider = (options: OauthBearerProviderOptions) => { const scope = options.scope && options.scope.trim().length > 0 ? options.scope.trim() : undefined; let tokenUrl: URL; try { - tokenUrl = new URL(options.host); + tokenUrl = new URL(options.endpointUrl); } catch { - throw new Error(`Kafka OAuth endpoint_url is not a valid URL: ${options.host}`); + throw new Error(`Kafka OAuth endpoint_url is not a valid URL: ${options.endpointUrl}`); } const client = new ClientCredentials({ client: { From 091c24ae6fa9e3193052d565b39563ebc2361d08 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 5 Jun 2026 11:09:22 +0300 Subject: [PATCH 12/15] Add Kafka OAUTHBEARER oauth config block to transport/vc-executor/edqs ymls for MSA mode --- docker/queue-confluent.env | 9 +++++++++ edqs/src/main/resources/edqs.yml | 12 ++++++++++++ .../src/main/resources/tb-vc-executor.yml | 12 ++++++++++++ .../coap/src/main/resources/tb-coap-transport.yml | 12 ++++++++++++ .../http/src/main/resources/tb-http-transport.yml | 12 ++++++++++++ .../lwm2m/src/main/resources/tb-lwm2m-transport.yml | 12 ++++++++++++ .../mqtt/src/main/resources/tb-mqtt-transport.yml | 12 ++++++++++++ .../snmp/src/main/resources/tb-snmp-transport.yml | 12 ++++++++++++ 8 files changed, 93 insertions(+) diff --git a/docker/queue-confluent.env b/docker/queue-confluent.env index 900504c4ea..616cb7316b 100644 --- a/docker/queue-confluent.env +++ b/docker/queue-confluent.env @@ -11,6 +11,15 @@ TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL=SASL_SSL TB_QUEUE_KAFKA_CONFLUENT_USERNAME=CLUSTER_API_KEY TB_QUEUE_KAFKA_CONFLUENT_PASSWORD=CLUSTER_API_SECRET +# Alternative: OAuth2/OIDC client-credentials (OAUTHBEARER) instead of PLAIN. +# Set TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM=OAUTHBEARER above and configure the values below. +# The token is fetched and refreshed automatically via Kafka's built-in OAuthBearerLoginCallbackHandler. +#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID= +#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET= +#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL= +# Optional: scope/audience required by some IdPs (e.g. Azure AD's "api:///.default") +#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_SCOPE= + TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 TB_QUEUE_KAFKA_TA_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 diff --git a/edqs/src/main/resources/edqs.yml b/edqs/src/main/resources/edqs.yml index cad0fcaccc..ac1b324eaa 100644 --- a/edqs/src/main/resources/edqs.yml +++ b/edqs/src/main/resources/edqs.yml @@ -134,11 +134,23 @@ 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:}" + # 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/msa/vc-executor/src/main/resources/tb-vc-executor.yml b/msa/vc-executor/src/main/resources/tb-vc-executor.yml index 024504f3c1..9d7a890690 100644 --- a/msa/vc-executor/src/main/resources/tb-vc-executor.yml +++ b/msa/vc-executor/src/main/resources/tb-vc-executor.yml @@ -107,11 +107,23 @@ 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:}" + # 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/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index 3bf2295ec9..c85662b478 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/transport/coap/src/main/resources/tb-coap-transport.yml @@ -323,11 +323,23 @@ 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:}" + # 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:}" # If you override any default Kafka topic name using environment variables, you must also specify the related consumer properties # for the new topic in `consumer-properties-per-topic-inline`. Otherwise, the topic will not inherit its expected configuration (e.g., max.poll.records, timeouts, etc). # Format: "topic1:key1=value1,key2=value2;topic2:key=value" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index 945a63240a..261edc70cf 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/transport/http/src/main/resources/tb-http-transport.yml @@ -272,11 +272,23 @@ queue: # The endpoint identification algorithm used by clients to validate server host name. 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 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:}" + # 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:}" # If you override any default Kafka topic name using environment variables, you must also specify the related consumer properties # for the new topic in `consumer-properties-per-topic-inline`. Otherwise, the topic will not inherit its expected configuration (e.g., max.poll.records, timeouts, etc). # Format: "topic1:key1=value1,key2=value2;topic2:key=value" diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 85068152a6..5f5fed8bae 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml @@ -373,11 +373,23 @@ 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:}" + # 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:}" # If you override any default Kafka topic name using environment variables, you must also specify the related consumer properties # for the new topic in `consumer-properties-per-topic-inline`. Otherwise, the topic will not inherit its expected configuration (e.g., max.poll.records, timeouts, etc). # Format: "topic1:key1=value1,key2=value2;topic2:key=value" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index 52d469e249..a193b1e4b6 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml @@ -306,11 +306,23 @@ 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:}" + # 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:}" # If you override any default Kafka topic name using environment variables, you must also specify the related consumer properties # for the new topic in `consumer-properties-per-topic-inline`. Otherwise, the topic will not inherit its expected configuration (e.g., max.poll.records, timeouts, etc). # Format: "topic1:key1=value1,key2=value2;topic2:key=value" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 46b1d8e48b..75d6b3f844 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -245,11 +245,23 @@ 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:}" + # 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:}" # If you override any default Kafka topic name using environment variables, you must also specify the related consumer properties # for the new topic in `consumer-properties-per-topic-inline`. Otherwise, the topic will not inherit its expected configuration (e.g., max.poll.records, timeouts, etc). # Format: "topic1:key1=value1,key2=value2;topic2:key=value" From 664b57457e819eaaca22b275ae0a940ea0e4d550 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 5 Jun 2026 11:11:21 +0300 Subject: [PATCH 13/15] Minor changes --- docker/queue-confluent.env | 9 --------- 1 file changed, 9 deletions(-) diff --git a/docker/queue-confluent.env b/docker/queue-confluent.env index 616cb7316b..900504c4ea 100644 --- a/docker/queue-confluent.env +++ b/docker/queue-confluent.env @@ -11,15 +11,6 @@ TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL=SASL_SSL TB_QUEUE_KAFKA_CONFLUENT_USERNAME=CLUSTER_API_KEY TB_QUEUE_KAFKA_CONFLUENT_PASSWORD=CLUSTER_API_SECRET -# Alternative: OAuth2/OIDC client-credentials (OAUTHBEARER) instead of PLAIN. -# Set TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM=OAUTHBEARER above and configure the values below. -# The token is fetched and refreshed automatically via Kafka's built-in OAuthBearerLoginCallbackHandler. -#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID= -#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET= -#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_ENDPOINT_URL= -# Optional: scope/audience required by some IdPs (e.g. Azure AD's "api:///.default") -#TB_QUEUE_KAFKA_CONFLUENT_OAUTH_SCOPE= - TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 TB_QUEUE_KAFKA_TA_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 From 89571c57c9f91c2e54f46994a70c2b7db0134ccb Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 8 Jun 2026 18:29:28 +0300 Subject: [PATCH 14/15] TbKafkaSettingsConfluentOAuthHttpWarningTest --- ...SettingsConfluentOAuthHttpWarningTest.java | 72 +++++++++++++++++++ .../config/custom-environment-variables.yml | 12 ++-- 2 files changed, 78 insertions(+), 6 deletions(-) create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java new file mode 100644 index 0000000000..15219da536 --- /dev/null +++ b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java @@ -0,0 +1,72 @@ +/** + * 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 ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; +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_PLAINTEXT", + "queue.kafka.confluent.oauth.client-id=my-client", + "queue.kafka.confluent.oauth.client-secret=my-secret", + "queue.kafka.confluent.oauth.endpoint-url=http://idp.example.com/oauth/token" +}) +class TbKafkaSettingsConfluentOAuthHttpWarningTest { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenHttpEndpointUrl_whenToProps_thenProducesValidJaasConfigAndLogsWarning() { + Logger logger = (Logger) LoggerFactory.getLogger(TbKafkaSettings.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + try { + 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\""); + assertThat(props).containsEntry("sasl.oauthbearer.token.endpoint.url", + "http://idp.example.com/oauth/token"); + + assertThat(appender.list) + .filteredOn(e -> e.getLevel() == Level.WARN) + .anySatisfy(e -> assertThat(e.getFormattedMessage()).contains("not HTTPS")); + } finally { + logger.detachAppender(appender); + } + } + +} diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 74c519b6fb..290ae4c219 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -54,17 +54,17 @@ kafka: mechanism: "TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM" username: "TB_QUEUE_KAFKA_CONFLUENT_USERNAME" password: "TB_QUEUE_KAFKA_CONFLUENT_PASSWORD" - # Optional: OAUTH Setting + # Optional: OAuth2 Setting oauth: - # Optional: Oauth Client ID + # Optional: OAuth2 Client ID client_id: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_ID" - # Optional: Oauth Client Secret + # Optional: OAuth2 Client Secret client_secret: "TB_QUEUE_KAFKA_CONFLUENT_OAUTH_CLIENT_SECRET" - # Optional: Oauth Endpoint URL to get the token from + # Optional: OAuth2 Endpoint URL to get the token from 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). Token refresh is managed by this service; no equivalent setting exists on the JVM-based services. 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") + # 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" logger: From 782af42ca14d4a3f3794c9cef1d06cac676b9483 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 8 Jun 2026 18:42:01 +0300 Subject: [PATCH 15/15] TbKafkaSettingsTest instead of multiples --- ...fkaSettingsConfluentOAuthEscapingTest.java | 56 ---- ...SettingsConfluentOAuthHttpWarningTest.java | 72 ---- ...ttingsConfluentOAuthMissingConfigTest.java | 48 --- ...bKafkaSettingsConfluentOAuthScopeTest.java | 56 ---- .../TbKafkaSettingsConfluentOAuthTest.java | 59 ---- .../TbKafkaSettingsConfluentPlainTest.java | 52 --- .../queue/kafka/TbKafkaSettingsTest.java | 308 +++++++++++++++--- msa/js-executor/config/default.yml | 2 + 8 files changed, 260 insertions(+), 393 deletions(-) delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java delete mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java 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 deleted file mode 100644 index 36857f8a2f..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthEscapingTest.java +++ /dev/null @@ -1,56 +0,0 @@ -/** - * 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=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/TbKafkaSettingsConfluentOAuthHttpWarningTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java deleted file mode 100644 index 15219da536..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthHttpWarningTest.java +++ /dev/null @@ -1,72 +0,0 @@ -/** - * 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 ch.qos.logback.classic.Level; -import ch.qos.logback.classic.Logger; -import ch.qos.logback.classic.spi.ILoggingEvent; -import ch.qos.logback.core.read.ListAppender; -import org.junit.jupiter.api.Test; -import org.slf4j.LoggerFactory; -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_PLAINTEXT", - "queue.kafka.confluent.oauth.client-id=my-client", - "queue.kafka.confluent.oauth.client-secret=my-secret", - "queue.kafka.confluent.oauth.endpoint-url=http://idp.example.com/oauth/token" -}) -class TbKafkaSettingsConfluentOAuthHttpWarningTest { - - @Autowired - TbKafkaSettings settings; - - @Test - void givenHttpEndpointUrl_whenToProps_thenProducesValidJaasConfigAndLogsWarning() { - Logger logger = (Logger) LoggerFactory.getLogger(TbKafkaSettings.class); - ListAppender appender = new ListAppender<>(); - appender.start(); - logger.addAppender(appender); - try { - 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\""); - assertThat(props).containsEntry("sasl.oauthbearer.token.endpoint.url", - "http://idp.example.com/oauth/token"); - - assertThat(appender.list) - .filteredOn(e -> e.getLevel() == Level.WARN) - .anySatisfy(e -> assertThat(e.getFormattedMessage()).contains("not HTTPS")); - } finally { - logger.detachAppender(appender); - } - } - -} 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 deleted file mode 100644 index 31f01f6121..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthMissingConfigTest.java +++ /dev/null @@ -1,48 +0,0 @@ -/** - * 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 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/TbKafkaSettingsConfluentOAuthScopeTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java deleted file mode 100644 index aa81d3131e..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthScopeTest.java +++ /dev/null @@ -1,56 +0,0 @@ -/** - * 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/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java deleted file mode 100644 index 4e56fd60fc..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentOAuthTest.java +++ /dev/null @@ -1,59 +0,0 @@ -/** - * 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 TbKafkaSettingsConfluentOAuthTest { - - @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/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 deleted file mode 100644 index 4414756610..0000000000 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsConfluentPlainTest.java +++ /dev/null @@ -1,52 +0,0 @@ -/** - * 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/TbKafkaSettingsTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsTest.java index 0175c45906..3792d76366 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/kafka/TbKafkaSettingsTest.java @@ -15,9 +15,15 @@ */ package org.thingsboard.server.queue.kafka; +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; import org.mockito.Mockito; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.TestPropertySource; @@ -25,76 +31,278 @@ import org.springframework.test.context.TestPropertySource; import java.util.Properties; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.spy; -@SpringBootTest(classes = {TbKafkaSettings.class, KafkaAdmin.class}) -@TestPropertySource(properties = { - "queue.type=kafka", - "queue.kafka.bootstrap.servers=localhost:9092", - "queue.kafka.other-inline=metrics.recording.level:INFO;metrics.sample.window.ms:30000", - "queue.kafka.consumer-properties-per-topic-inline=" + - "tb_core_updated:max.poll.records=10;" + - "tb_core_updated:enable.auto.commit=true;" + - "tb_core_updated:bootstrap.servers=kafka1:9092,kafka2:9092;" + - "tb_edge_updated:max.poll.records=5;" + - "tb_edge_updated:auto.offset.reset=latest" -}) class TbKafkaSettingsTest { - @Autowired - TbKafkaSettings settings; + @Nested + @SpringBootTest(classes = {TbKafkaSettings.class, KafkaAdmin.class}) + @TestPropertySource(properties = { + "queue.type=kafka", + "queue.kafka.bootstrap.servers=localhost:9092", + "queue.kafka.other-inline=metrics.recording.level:INFO;metrics.sample.window.ms:30000", + "queue.kafka.consumer-properties-per-topic-inline=" + + "tb_core_updated:max.poll.records=10;" + + "tb_core_updated:enable.auto.commit=true;" + + "tb_core_updated:bootstrap.servers=kafka1:9092,kafka2:9092;" + + "tb_edge_updated:max.poll.records=5;" + + "tb_edge_updated:auto.offset.reset=latest" + }) + class InlinePropertiesAndSsl { + + @Autowired + TbKafkaSettings settings; + + @BeforeEach + void beforeEach() { + settings = spy(settings); // SpyBean is not aware on @ConditionalOnProperty, that is why the traditional spy in use + } + + @Test + void givenToProps_whenConfigureSSL_thenVerifyOnce() { + Properties props = settings.toProps(); + + assertThat(props).as("TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS").containsEntry("request.timeout.ms", 30000); + + //other-inline + assertThat(props).as("metrics.recording.level").containsEntry("metrics.recording.level", "INFO"); + assertThat(props).as("TB_QUEUE_KAFKA_SESSION_TIMEOUT_MS").containsEntry("metrics.sample.window.ms", "30000"); + + Mockito.verify(settings).toProps(); + Mockito.verify(settings).configureSSL(any()); + } + + @Test + void givenToAdminProps_whenConfigureSSL_thenVerifyOnce() { + settings.toAdminProps(); + Mockito.verify(settings).toProps(); + Mockito.verify(settings).configureSSL(any()); + } + + @Test + void givenToConsumerProps_whenConfigureSSL_thenVerifyOnce() { + settings.toConsumerProps("main"); + Mockito.verify(settings).toProps(); + Mockito.verify(settings).configureSSL(any()); + } + + @Test + void givenTotoProducerProps_whenConfigureSSL_thenVerifyOnce() { + settings.toProducerProps(); + Mockito.verify(settings).toProps(); + Mockito.verify(settings).configureSSL(any()); + } + + @Test + void givenMultipleTopicsInInlineConfig_whenParsed_thenEachTopicGetsExpectedProperties() { + Properties coreProps = settings.toConsumerProps("tb_core_updated"); + assertThat(coreProps.getProperty("max.poll.records")).isEqualTo("10"); + assertThat(coreProps.getProperty("enable.auto.commit")).isEqualTo("true"); + assertThat(coreProps.getProperty("bootstrap.servers")).isEqualTo("kafka1:9092,kafka2:9092"); + + Properties edgeProps = settings.toConsumerProps("tb_edge_updated"); + assertThat(edgeProps.getProperty("max.poll.records")).isEqualTo("5"); + assertThat(edgeProps.getProperty("auto.offset.reset")).isEqualTo("latest"); + } - @BeforeEach - void beforeEach() { - settings = spy(settings); //SpyBean is not aware on @ConditionalOnProperty, that is why the traditional spy in use } - @Test - void givenToProps_whenConfigureSSL_thenVerifyOnce() { - Properties props = settings.toProps(); + @Nested + @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 ConfluentOAuth { - assertThat(props).as("TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS").containsEntry("request.timeout.ms", 30000); + @Autowired + TbKafkaSettings settings; - //other-inline - assertThat(props).as("metrics.recording.level").containsEntry("metrics.recording.level", "INFO"); - assertThat(props).as("TB_QUEUE_KAFKA_SESSION_TIMEOUT_MS").containsEntry("metrics.sample.window.ms", "30000"); + @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\""); + } - Mockito.verify(settings).toProps(); - Mockito.verify(settings).configureSSL(any()); } - @Test - void givenToAdminProps_whenConfigureSSL_thenVerifyOnce() { - settings.toAdminProps(); - Mockito.verify(settings).toProps(); - Mockito.verify(settings).configureSSL(any()); + @Nested + @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 ConfluentOAuthEscaping { + + @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("\";"); + } + } - @Test - void givenToConsumerProps_whenConfigureSSL_thenVerifyOnce() { - settings.toConsumerProps("main"); - Mockito.verify(settings).toProps(); - Mockito.verify(settings).configureSSL(any()); + @Nested + @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_PLAINTEXT", + "queue.kafka.confluent.oauth.client-id=my-client", + "queue.kafka.confluent.oauth.client-secret=my-secret", + "queue.kafka.confluent.oauth.endpoint-url=http://idp.example.com/oauth/token" + }) + class ConfluentOAuthHttpWarning { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenHttpEndpointUrl_whenToProps_thenProducesValidJaasConfigAndLogsWarning() { + Logger logger = (Logger) LoggerFactory.getLogger(TbKafkaSettings.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + try { + 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\""); + assertThat(props).containsEntry("sasl.oauthbearer.token.endpoint.url", + "http://idp.example.com/oauth/token"); + + assertThat(appender.list) + .filteredOn(e -> e.getLevel() == Level.WARN) + .anySatisfy(e -> assertThat(e.getFormattedMessage()).contains("not HTTPS")); + } finally { + logger.detachAppender(appender); + } + } + + } + + @Nested + @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 ConfluentOAuthScope { + + @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("\";"); + } + } - @Test - void givenTotoProducerProps_whenConfigureSSL_thenVerifyOnce() { - settings.toProducerProps(); - Mockito.verify(settings).toProps(); - Mockito.verify(settings).configureSSL(any()); + @Nested + @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 ConfluentOAuthMissingConfig { + + @Autowired + TbKafkaSettings settings; + + @Test + void givenOauthbearerWithoutEndpointUrl_whenToProps_thenFailsFast() { + assertThatThrownBy(() -> settings.toProps()) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("OAUTHBEARER") + .hasMessageContaining("endpoint-url"); + } + } - @Test - void givenMultipleTopicsInInlineConfig_whenParsed_thenEachTopicGetsExpectedProperties() { - Properties coreProps = settings.toConsumerProps("tb_core_updated"); - assertThat(coreProps.getProperty("max.poll.records")).isEqualTo("10"); - assertThat(coreProps.getProperty("enable.auto.commit")).isEqualTo("true"); - assertThat(coreProps.getProperty("bootstrap.servers")).isEqualTo("kafka1:9092,kafka2:9092"); + @Nested + @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 ConfluentPlain { + + @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"); + } - Properties edgeProps = settings.toConsumerProps("tb_edge_updated"); - assertThat(edgeProps.getProperty("max.poll.records")).isEqualTo("5"); - assertThat(edgeProps.getProperty("auto.offset.reset")).isEqualTo("latest"); } } diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 8acaaaaba9..38360f9e8b 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -49,6 +49,8 @@ kafka: confluent: sasl: mechanism: "PLAIN" # SASL mechanism for Confluent Cloud authentication: PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512 + oauth: + refresh_threshold: "60000" # How many ms before token expiry to proactively refresh it (default: 60 s) # Logging configuration # Controls log verbosity, output directory, and log file naming pattern for the JS executor service.