From 37680608e139852cbb68defc0d372f6622ef9d7c Mon Sep 17 00:00:00 2001 From: Shizuo Fujita Date: Thu, 3 Sep 2026 17:41:10 +0900 Subject: [PATCH 1/2] in_rdkafka_group: Apply SSL and SASL parameters to the consumer in_rdkafka_group accepted the ssl_* and SASL parameters but never passed them to librdkafka. The consumer was built from kafka_configs alone, so a configuration that only set ssl_ca_cert connected over PLAINTEXT without any error. Move the parameter mapping of out_rdkafka2#build_config into a shared module and use it from both plugins. Keys given in kafka_configs take precedence, matching rdkafka_options in out_rdkafka2. The sasl_over_ssl check now applies to in_rdkafka_group as well, and it matches security.protocol case-insensitively so that the lowercase sasl_plaintext spelling used in the librdkafka documentation cannot bypass it. Signed-off-by: Shizuo Fujita Co-Authored-By: Claude Fable 5.1 --- README.md | 6 + lib/fluent/plugin/in_rdkafka_group.rb | 14 ++- lib/fluent/plugin/kafka_plugin_util.rb | 64 ++++++++++ lib/fluent/plugin/out_rdkafka2.rb | 51 +------- test/plugin/test_in_rdkafka_group.rb | 158 +++++++++++++++++++++++++ test/plugin/test_out_rdkafka2.rb | 9 ++ 6 files changed, 253 insertions(+), 49 deletions(-) diff --git a/README.md b/README.md index 9ca58a53..138eac2c 100644 --- a/README.md +++ b/README.md @@ -166,6 +166,12 @@ With the introduction of the rdkafka-ruby based input plugin we hope to support time_source :default => now time_format + # The SSL and SASL parameters in "Common parameters" are mapped to librdkafka + # properties. Keys given in kafka_configs take precedence over the mapped values. + # ssl_ca_cert uses only its first entry, and ssl_client_cert_chain is not mapped. + # When true, SASL authentication requires an SSL connection + sasl_over_ssl (bool) :default => true + # kafka consumer options max_wait_time_ms 500 max_batch_size 10000 diff --git a/lib/fluent/plugin/in_rdkafka_group.rb b/lib/fluent/plugin/in_rdkafka_group.rb index 3f85d53d..20957452 100644 --- a/lib/fluent/plugin/in_rdkafka_group.rb +++ b/lib/fluent/plugin/in_rdkafka_group.rb @@ -56,6 +56,7 @@ class Fluent::Plugin::RdKafkaGroupInput < Fluent::Plugin::Input include Fluent::KafkaPluginUtil::SSLSettings include Fluent::KafkaPluginUtil::SaslSettings + include Fluent::KafkaPluginUtil::RdkafkaSecuritySettings class ForShutdown < StandardError end @@ -126,6 +127,8 @@ def configure(conf) if @time_source == :record and @time_format @time_parser = Fluent::TimeParser.new(@time_format) end + + @rdkafka_config = build_config end def setup_parser(parser_conf) @@ -177,8 +180,17 @@ def shutdown super end + def build_config + config = rdkafka_security_config + @kafka_configs.each { |k, v| + config[k.to_sym] = v + } + validate_sasl_over_ssl(config) + config + end + def setup_consumer - consumer = Rdkafka::Config.new(@kafka_configs).consumer + consumer = Rdkafka::Config.new(@rdkafka_config).consumer consumer.subscribe(*@topics) consumer end diff --git a/lib/fluent/plugin/kafka_plugin_util.rb b/lib/fluent/plugin/kafka_plugin_util.rb index a8665adc..fd825a12 100644 --- a/lib/fluent/plugin/kafka_plugin_util.rb +++ b/lib/fluent/plugin/kafka_plugin_util.rb @@ -119,5 +119,69 @@ def configure(conf) @scram_mechanism = @scram_mechanism.to_s if @scram_mechanism end end + + module RdkafkaSecuritySettings + SCRAM_MECHANISMS = { + "sha256" => "SCRAM-SHA-256", + "sha512" => "SCRAM-SHA-512", + }.freeze + + def self.included(klass) + klass.instance_eval { + config_param :service_name, :string, :default => nil, :desc => 'Used for sasl.kerberos.service.name' + config_param :sasl_over_ssl, :bool, :default => true, + :desc => 'When true, SASL authentication requires an SSL connection' + } + end + + def rdkafka_security_config + config = {} + + if (@ssl_ca_cert && @ssl_ca_cert[0]) || @ssl_ca_certs_from_system || @ssl_client_cert || @ssl_client_cert_key + ssl = true + config[:"ssl.ca.location"] = @ssl_ca_cert[0] if @ssl_ca_cert && @ssl_ca_cert[0] + config[:"ssl.certificate.location"] = @ssl_client_cert if @ssl_client_cert + config[:"ssl.key.location"] = @ssl_client_cert_key if @ssl_client_cert_key + config[:"ssl.key.password"] = @ssl_client_cert_key_password if @ssl_client_cert_key_password + config[:"ssl.endpoint.identification.algorithm"] = @ssl_verify_hostname ? "https" : "none" + config[:"enable.ssl.certificate.verification"] = true + end + + if @principal + sasl = true + config[:"sasl.mechanisms"] = "GSSAPI" + config[:"sasl.kerberos.principal"] = @principal + config[:"sasl.kerberos.service.name"] = @service_name if @service_name + config[:"sasl.kerberos.keytab"] = @keytab if @keytab + end + + if @username && @password + sasl = true + config[:"sasl.mechanisms"] = SCRAM_MECHANISMS.fetch(@scram_mechanism, 'PLAIN') + end + + if ssl && sasl + security_protocol = "SASL_SSL" + elsif ssl && !sasl + security_protocol = "SSL" + elsif !ssl && sasl + security_protocol = "SASL_PLAINTEXT" + else + security_protocol = "PLAINTEXT" + end + config[:"security.protocol"] = security_protocol + + config[:"sasl.username"] = @username if @username + config[:"sasl.password"] = @password if @password + + config + end + + def validate_sasl_over_ssl(config) + if @sasl_over_ssl && config[:"sasl.password"] && config[:"security.protocol"].to_s.upcase == "SASL_PLAINTEXT" + raise Fluent::ConfigError, "SASL authentication requires that SSL is configured. Set 'sasl_over_ssl false' to send SASL credentials over a plaintext connection" + end + end + end end end diff --git a/lib/fluent/plugin/out_rdkafka2.rb b/lib/fluent/plugin/out_rdkafka2.rb index 4b0c4698..fc45a063 100644 --- a/lib/fluent/plugin/out_rdkafka2.rb +++ b/lib/fluent/plugin/out_rdkafka2.rb @@ -99,9 +99,6 @@ class Fluent::Rdkafka2Output < Output config_param :enqueue_retry_backoff, :integer, :default => 3 config_param :max_enqueue_bytes_per_second, :size, :default => nil, :desc => 'The maximum number of enqueueing bytes per second' - config_param :service_name, :string, :default => nil, :desc => 'Used for sasl.kerberos.service.name' - config_param :sasl_over_ssl, :bool, :default => true, - :desc => 'When true, SASL authentication requires an SSL connection' config_param :unrecoverable_error_codes, :array, :default => ["topic_authorization_failed", "msg_size_too_large"], :desc => 'Handle some of the error codes should be unrecoverable if specified' @@ -116,11 +113,7 @@ class Fluent::Rdkafka2Output < Output include Fluent::KafkaPluginUtil::SSLSettings include Fluent::KafkaPluginUtil::SaslSettings include Fluent::KafkaPluginUtil::PartitionSettings - - SCRAM_MECHANISMS = { - "sha256" => "SCRAM-SHA-256", - "sha512" => "SCRAM-SHA-512", - }.freeze + include Fluent::KafkaPluginUtil::RdkafkaSecuritySettings class EnqueueRate class LimitExceeded < StandardError @@ -260,41 +253,7 @@ def add(level, message = nil) end def build_config - config = {:"bootstrap.servers" => @brokers} - - if (@ssl_ca_cert && @ssl_ca_cert[0]) || @ssl_ca_certs_from_system || @ssl_client_cert || @ssl_client_cert_key - ssl = true - config[:"ssl.ca.location"] = @ssl_ca_cert[0] if @ssl_ca_cert && @ssl_ca_cert[0] - config[:"ssl.certificate.location"] = @ssl_client_cert if @ssl_client_cert - config[:"ssl.key.location"] = @ssl_client_cert_key if @ssl_client_cert_key - config[:"ssl.key.password"] = @ssl_client_cert_key_password if @ssl_client_cert_key_password - config[:"ssl.endpoint.identification.algorithm"] = @ssl_verify_hostname ? "https" : "none" - config[:"enable.ssl.certificate.verification"] = true - end - - if @principal - sasl = true - config[:"sasl.mechanisms"] = "GSSAPI" - config[:"sasl.kerberos.principal"] = @principal - config[:"sasl.kerberos.service.name"] = @service_name if @service_name - config[:"sasl.kerberos.keytab"] = @keytab if @keytab - end - - if @username && @password - sasl = true - config[:"sasl.mechanisms"] = SCRAM_MECHANISMS.fetch(@scram_mechanism, 'PLAIN') - end - - if ssl && sasl - security_protocol = "SASL_SSL" - elsif ssl && !sasl - security_protocol = "SSL" - elsif !ssl && sasl - security_protocol = "SASL_PLAINTEXT" - else - security_protocol = "PLAINTEXT" - end - config[:"security.protocol"] = security_protocol + config = {:"bootstrap.servers" => @brokers}.merge(rdkafka_security_config) config[:"compression.codec"] = @compression_codec if @compression_codec config[:"message.send.max.retries"] = @max_send_retries if @max_send_retries @@ -304,17 +263,13 @@ def build_config config[:"queue.buffering.max.messages"] = @rdkafka_buffering_max_messages if @rdkafka_buffering_max_messages config[:"message.max.bytes"] = @rdkafka_message_max_bytes if @rdkafka_message_max_bytes config[:"batch.num.messages"] = @rdkafka_message_max_num if @rdkafka_message_max_num - config[:"sasl.username"] = @username if @username - config[:"sasl.password"] = @password if @password config[:"enable.idempotence"] = @idempotent if @idempotent @rdkafka_options.each { |k, v| config[k.to_sym] = v } - if @sasl_over_ssl && config[:"sasl.password"] && config[:"security.protocol"] == "SASL_PLAINTEXT" - raise Fluent::ConfigError, "SASL authentication requires that SSL is configured. Set 'sasl_over_ssl false' to send SASL credentials over a plaintext connection" - end + validate_sasl_over_ssl(config) config end diff --git a/test/plugin/test_in_rdkafka_group.rb b/test/plugin/test_in_rdkafka_group.rb index 2ab92ea0..af6d7fe6 100644 --- a/test/plugin/test_in_rdkafka_group.rb +++ b/test/plugin/test_in_rdkafka_group.rb @@ -44,6 +44,164 @@ def test_multi_worker_support assert_true d.instance.multi_workers_ready? end + def test_build_config_without_security_parameters + d = create_driver + config = d.instance.build_config + + assert_equal 'localhost:9092', config[:"bootstrap.servers"] + assert_equal 'test_group', config[:"group.id"] + assert_equal 'PLAINTEXT', config[:"security.protocol"] + assert_nil config[:"ssl.endpoint.identification.algorithm"] + assert_nil config[:"sasl.mechanisms"] + end + + def test_build_config_ssl_ca_certs_from_system + d = create_driver(CONFIG + "ssl_ca_certs_from_system true\n") + config = d.instance.build_config + + assert_equal 'SSL', config[:"security.protocol"] + assert_equal 'https', config[:"ssl.endpoint.identification.algorithm"] + assert_equal true, config[:"enable.ssl.certificate.verification"] + assert_nil config[:"ssl.ca.location"] + end + + def test_build_config_ssl_client_cert + d = create_driver(CONFIG + %[ + ssl_ca_cert /path/to/ca_cert.pem + ssl_client_cert /path/to/cert.pem + ssl_client_cert_key /path/to/key.pem + ssl_client_cert_key_password secret + ]) + config = d.instance.build_config + + assert_equal 'SSL', config[:"security.protocol"] + assert_equal '/path/to/ca_cert.pem', config[:"ssl.ca.location"] + assert_equal '/path/to/cert.pem', config[:"ssl.certificate.location"] + assert_equal '/path/to/key.pem', config[:"ssl.key.location"] + assert_equal 'secret', config[:"ssl.key.password"] + end + + def test_build_config_ssl_verify_hostname_false + d = create_driver(CONFIG + "ssl_ca_certs_from_system true\nssl_verify_hostname false\n") + config = d.instance.build_config + + assert_equal 'none', config[:"ssl.endpoint.identification.algorithm"] + assert_equal true, config[:"enable.ssl.certificate.verification"] + end + + def test_build_config_sasl_plain_over_ssl + d = create_driver(CONFIG + "username testuser\npassword testpass\nssl_ca_certs_from_system true\n") + config = d.instance.build_config + + assert_equal 'SASL_SSL', config[:"security.protocol"] + assert_equal 'PLAIN', config[:"sasl.mechanisms"] + assert_equal 'testuser', config[:"sasl.username"] + assert_equal 'testpass', config[:"sasl.password"] + end + + def test_configure_sasl_plain_without_ssl_raises + assert_raise(Fluent::ConfigError) { + create_driver(CONFIG + "username testuser\npassword testpass\n") + } + end + + def test_build_config_sasl_plain_without_ssl_allowed_by_sasl_over_ssl + d = create_driver(CONFIG + "username testuser\npassword testpass\nsasl_over_ssl false\n") + config = d.instance.build_config + + assert_equal 'SASL_PLAINTEXT', config[:"security.protocol"] + assert_equal 'testpass', config[:"sasl.password"] + end + + data("sha256" => ["sha256", "SCRAM-SHA-256"], + "sha512" => ["sha512", "SCRAM-SHA-512"]) + def test_build_config_sasl_scram(data) + mechanism, expected = data + d = create_driver(CONFIG + "username testuser\npassword testpass\nscram_mechanism #{mechanism}\nssl_ca_certs_from_system true\n") + config = d.instance.build_config + + assert_equal expected, config[:"sasl.mechanisms"] + assert_equal 'SASL_SSL', config[:"security.protocol"] + end + + def test_build_config_sasl_gssapi + d = create_driver(CONFIG + "principal kafka/host@REALM\nkeytab /path/to/kafka.keytab\nservice_name kafka\nssl_ca_certs_from_system true\n") + config = d.instance.build_config + + assert_equal 'SASL_SSL', config[:"security.protocol"] + assert_equal 'GSSAPI', config[:"sasl.mechanisms"] + assert_equal 'kafka/host@REALM', config[:"sasl.kerberos.principal"] + assert_equal '/path/to/kafka.keytab', config[:"sasl.kerberos.keytab"] + assert_equal 'kafka', config[:"sasl.kerberos.service.name"] + end + + def test_build_config_kafka_configs_take_precedence + conf = %[ + topics #{TOPIC_NAME} + ssl_ca_certs_from_system true + kafka_configs {"bootstrap.servers": "localhost:9092", "group.id": "test_group", "security.protocol": "SASL_SSL", "ssl.endpoint.identification.algorithm": "none"} + + @type none + + ] + d = create_driver(conf) + config = d.instance.build_config + + assert_equal 'SASL_SSL', config[:"security.protocol"] + assert_equal 'none', config[:"ssl.endpoint.identification.algorithm"] + assert_nil config["security.protocol"] + end + + def test_configure_sasl_plaintext_in_kafka_configs_raises + conf = %[ + topics #{TOPIC_NAME} + kafka_configs {"bootstrap.servers": "localhost:9092", "group.id": "test_group", "security.protocol": "SASL_PLAINTEXT", "sasl.mechanisms": "PLAIN", "sasl.username": "testuser", "sasl.password": "testpass"} + + @type none + + ] + + assert_raise(Fluent::ConfigError) { + create_driver(conf) + } + assert_nothing_raised { + create_driver(conf + "sasl_over_ssl false\n") + } + end + + data("uppercase" => "SASL_PLAINTEXT", + "lowercase" => "sasl_plaintext") + def test_configure_sasl_plaintext_in_kafka_configs_raises_regardless_of_case(protocol) + conf = %[ + topics #{TOPIC_NAME} + username testuser + password testpass + kafka_configs {"bootstrap.servers": "localhost:9092", "group.id": "test_group", "security.protocol": "#{protocol}"} + + @type none + + ] + + assert_raise(Fluent::ConfigError) { + create_driver(conf) + } + end + + def test_setup_consumer_uses_build_config + d = create_driver(CONFIG + "ssl_ca_certs_from_system true\n") + consumer = Object.new + stub(consumer).subscribe + rdkafka_config = Object.new + stub(rdkafka_config).consumer { consumer } + passed = nil + stub(Rdkafka::Config).new { |config| passed = config; rdkafka_config } + + d.instance.setup_consumer + + assert_equal 'localhost:9092', passed[:"bootstrap.servers"] + assert_equal 'SSL', passed[:"security.protocol"] + end + class ConsumeTest < self TOPIC_NAME = "kafka-input-#{SecureRandom.uuid}" diff --git a/test/plugin/test_out_rdkafka2.rb b/test/plugin/test_out_rdkafka2.rb index 017bf2ae..bdeeec98 100644 --- a/test/plugin/test_out_rdkafka2.rb +++ b/test/plugin/test_out_rdkafka2.rb @@ -173,6 +173,15 @@ def test_configure_sasl_plain_with_security_protocol_from_rdkafka_options assert_equal 'SASL_SSL', config[:"security.protocol"] end + def test_configure_sasl_plain_with_lowercase_security_protocol_from_rdkafka_options + conf = base_config + config_element('ROOT', '', {"username" => "testuser", "password" => "testpass", + "rdkafka_options" => '{"security.protocol": "sasl_plaintext"}'}, []) + + assert_raise(Fluent::ConfigError) { + create_driver(conf) + } + end + data("sha256" => ["sha256", "SCRAM-SHA-256"], "sha512" => ["sha512", "SCRAM-SHA-512"]) def test_configure_sasl_scram(data) From fd971de974447a2cc5f742660d2134f83a8de33c Mon Sep 17 00:00:00 2001 From: Shizuo Fujita Date: Fri, 4 Sep 2026 10:49:36 +0900 Subject: [PATCH 2/2] in_rdkafka_group, out_rdkafka2: Warn when scram_mechanism has no credentials scram_mechanism only takes effect when username and password are set, so a configuration that sets it alone was ignored without any message. Also switch the test configurations of in_rdkafka_group to the %[] style already used in that file. Co-Authored-By: Claude Opus 5 (1M context) Signed-off-by: Shizuo Fujita --- lib/fluent/plugin/kafka_plugin_util.rb | 2 + test/plugin/test_in_rdkafka_group.rb | 56 ++++++++++++++++++++++---- 2 files changed, 50 insertions(+), 8 deletions(-) diff --git a/lib/fluent/plugin/kafka_plugin_util.rb b/lib/fluent/plugin/kafka_plugin_util.rb index fd825a12..949ce7d8 100644 --- a/lib/fluent/plugin/kafka_plugin_util.rb +++ b/lib/fluent/plugin/kafka_plugin_util.rb @@ -158,6 +158,8 @@ def rdkafka_security_config if @username && @password sasl = true config[:"sasl.mechanisms"] = SCRAM_MECHANISMS.fetch(@scram_mechanism, 'PLAIN') + elsif @scram_mechanism + log.warn "scram_mechanism is ignored because username and password are not set" end if ssl && sasl diff --git a/test/plugin/test_in_rdkafka_group.rb b/test/plugin/test_in_rdkafka_group.rb index af6d7fe6..b37ada47 100644 --- a/test/plugin/test_in_rdkafka_group.rb +++ b/test/plugin/test_in_rdkafka_group.rb @@ -56,7 +56,9 @@ def test_build_config_without_security_parameters end def test_build_config_ssl_ca_certs_from_system - d = create_driver(CONFIG + "ssl_ca_certs_from_system true\n") + d = create_driver(CONFIG + %[ + ssl_ca_certs_from_system true + ]) config = d.instance.build_config assert_equal 'SSL', config[:"security.protocol"] @@ -82,7 +84,10 @@ def test_build_config_ssl_client_cert end def test_build_config_ssl_verify_hostname_false - d = create_driver(CONFIG + "ssl_ca_certs_from_system true\nssl_verify_hostname false\n") + d = create_driver(CONFIG + %[ + ssl_ca_certs_from_system true + ssl_verify_hostname false + ]) config = d.instance.build_config assert_equal 'none', config[:"ssl.endpoint.identification.algorithm"] @@ -90,7 +95,11 @@ def test_build_config_ssl_verify_hostname_false end def test_build_config_sasl_plain_over_ssl - d = create_driver(CONFIG + "username testuser\npassword testpass\nssl_ca_certs_from_system true\n") + d = create_driver(CONFIG + %[ + username testuser + password testpass + ssl_ca_certs_from_system true + ]) config = d.instance.build_config assert_equal 'SASL_SSL', config[:"security.protocol"] @@ -101,12 +110,19 @@ def test_build_config_sasl_plain_over_ssl def test_configure_sasl_plain_without_ssl_raises assert_raise(Fluent::ConfigError) { - create_driver(CONFIG + "username testuser\npassword testpass\n") + create_driver(CONFIG + %[ + username testuser + password testpass + ]) } end def test_build_config_sasl_plain_without_ssl_allowed_by_sasl_over_ssl - d = create_driver(CONFIG + "username testuser\npassword testpass\nsasl_over_ssl false\n") + d = create_driver(CONFIG + %[ + username testuser + password testpass + sasl_over_ssl false + ]) config = d.instance.build_config assert_equal 'SASL_PLAINTEXT', config[:"security.protocol"] @@ -117,15 +133,37 @@ def test_build_config_sasl_plain_without_ssl_allowed_by_sasl_over_ssl "sha512" => ["sha512", "SCRAM-SHA-512"]) def test_build_config_sasl_scram(data) mechanism, expected = data - d = create_driver(CONFIG + "username testuser\npassword testpass\nscram_mechanism #{mechanism}\nssl_ca_certs_from_system true\n") + d = create_driver(CONFIG + %[ + username testuser + password testpass + scram_mechanism #{mechanism} + ssl_ca_certs_from_system true + ]) config = d.instance.build_config assert_equal expected, config[:"sasl.mechanisms"] assert_equal 'SASL_SSL', config[:"security.protocol"] end + def test_build_config_sasl_scram_without_credentials_warns + d = create_driver(CONFIG + %[ + scram_mechanism sha256 + ssl_ca_certs_from_system true + ]) + config = d.instance.build_config + + assert_equal 'SSL', config[:"security.protocol"] + assert_nil config[:"sasl.mechanisms"] + assert_true d.logs.any? { |log| log.include?("scram_mechanism is ignored") } + end + def test_build_config_sasl_gssapi - d = create_driver(CONFIG + "principal kafka/host@REALM\nkeytab /path/to/kafka.keytab\nservice_name kafka\nssl_ca_certs_from_system true\n") + d = create_driver(CONFIG + %[ + principal kafka/host@REALM + keytab /path/to/kafka.keytab + service_name kafka + ssl_ca_certs_from_system true + ]) config = d.instance.build_config assert_equal 'SASL_SSL', config[:"security.protocol"] @@ -188,7 +226,9 @@ def test_configure_sasl_plaintext_in_kafka_configs_raises_regardless_of_case(pro end def test_setup_consumer_uses_build_config - d = create_driver(CONFIG + "ssl_ca_certs_from_system true\n") + d = create_driver(CONFIG + %[ + ssl_ca_certs_from_system true + ]) consumer = Object.new stub(consumer).subscribe rdkafka_config = Object.new