diff --git a/README.md b/README.md index 9ca58a5..138eac2 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 3f85d53..2095745 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 a8665ad..949ce7d 100644 --- a/lib/fluent/plugin/kafka_plugin_util.rb +++ b/lib/fluent/plugin/kafka_plugin_util.rb @@ -119,5 +119,71 @@ 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') + elsif @scram_mechanism + log.warn "scram_mechanism is ignored because username and password are not set" + 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 4b0c469..fc45a06 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 2ab92ea..b37ada4 100644 --- a/test/plugin/test_in_rdkafka_group.rb +++ b/test/plugin/test_in_rdkafka_group.rb @@ -44,6 +44,204 @@ 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 + ]) + 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 + ssl_verify_hostname false + ]) + 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 + password testpass + ssl_ca_certs_from_system true + ]) + 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 + password testpass + ]) + } + end + + def test_build_config_sasl_plain_without_ssl_allowed_by_sasl_over_ssl + d = create_driver(CONFIG + %[ + username testuser + password testpass + sasl_over_ssl false + ]) + 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 + 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 + 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"] + 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 + ]) + 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 017bf2a..bdeeec9 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)