diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6ceef7bc2e..ca450417dd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1782,11 +1782,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/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..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 @@ -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; @@ -32,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; @@ -137,6 +139,18 @@ 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.confluent.oauth.scope:}") + private String oauthScope; + @Value("${queue.kafka.other-inline:}") private String otherInline; @@ -213,9 +227,13 @@ 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)) { + applyOauthBearerProps(props); + } else { + props.put(SaslConfigs.SASL_JAAS_CONFIG, saslConfig); + } } props.put(CommonClientConfigs.REQUEST_TIMEOUT_MS_CONFIG, requestTimeoutMs); @@ -230,6 +248,34 @@ 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.regionMatches(true, 0, "https://", 0, "https://".length())) { + log.warn("Kafka OAuth token endpoint URL is not HTTPS ({}); client credentials will be sent unencrypted", + oauthEndpointUrl); + } + StringBuilder jaasConfig = new StringBuilder( + "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required" + + " clientId=\"" + escapeJaasValue(oauthClientId) + "\"" + + " 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); + } + + 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/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/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/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 9923fe4ed3..290ae4c219 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -54,6 +54,18 @@ kafka: mechanism: "TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM" username: "TB_QUEUE_KAFKA_CONFLUENT_USERNAME" password: "TB_QUEUE_KAFKA_CONFLUENT_PASSWORD" + # Optional: OAuth2 Setting + 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 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). 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: 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: level: "LOGGER_LEVEL" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index feae8edc4a..c84ae1db25 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. diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 6aff900f30..75df2f40f4 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -19,6 +19,7 @@ "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" @@ -35,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/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index 175abd80f4..95c95c6e39 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, CompressionCodecs, @@ -36,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`); @@ -114,11 +117,35 @@ 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') - }; + 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; + 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: optionalOauthStr('kafka.confluent.oauth.client_id'), + clientSecret: optionalOauthStr('kafka.confluent.oauth.client_secret'), + endpointUrl: optionalOauthStr('kafka.confluent.oauth.endpoint_url'), + refreshThresholdMs, + scope, + }) + }; + } else { + kafkaConfig['sasl'] = { + mechanism: saslMechanism 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..83ae5125b6 --- /dev/null +++ b/msa/js-executor/queue/oAuthBearerProvider.ts @@ -0,0 +1,171 @@ +/// +/// 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'; +import { _logger } from '../config/logger'; + +interface OauthBearerProviderOptions { + clientId: string; + clientSecret: string; + endpointUrl: 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; +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'); + 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.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) { + 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.endpointUrl); + } catch { + throw new Error(`Kafka OAuth endpoint_url is not a valid URL: ${options.endpointUrl}`); + } + const client = new ClientCredentials({ + client: { + id: options.clientId, + secret: options.clientSecret + }, + auth: { + // 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 + } + }); + + // 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; + // 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; + + 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(() => { + // refreshTokenOnce() updates cachedToken on success and reschedules its own retry on failure; + // swallow here so a failed background refresh is not an unhandled rejection. + 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 { + // Client-credentials grant issues no refresh_token, so always request a fresh token. + const accessToken: AccessToken = await client.getToken(scope ? { scope } : {}); + + // 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 lifetimeMs = expiresIn * 1000; + cachedToken = accessTokenValue; + cachedTokenExpiresAt = Date.now() + lifetimeMs; + + // 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) { + logger.error('Failed to obtain Kafka OAuth bearer token, retrying in %dms', RETRY_DELAY_MS); + scheduleRefresh(RETRY_DELAY_MS); + throw err; + } + } + + // 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. + refreshTokenOnce().catch(() => { /* retry already scheduled in refreshToken() */ }); + + return async function () { + // 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 + }; + }; +}; diff --git a/msa/js-executor/yarn.lock b/msa/js-executor/yarn.lock index 5e3330e0d7..c11a29f1b7 100644 --- a/msa/js-executor/yarn.lock +++ b/msa/js-executor/yarn.lock @@ -131,6 +131,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.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" + "@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" @@ -184,6 +222,23 @@ dependencies: "@napi-rs/triples" "^1.2.0" +"@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" @@ -294,6 +349,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" @@ -623,6 +683,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" @@ -1029,6 +1096,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" @@ -1549,6 +1627,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" 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"