Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,12 @@ With the introduction of the rdkafka-ruby based input plugin we hope to support
time_source <source for message timestamp (now|kafka|record)> :default => now
time_format <string (Optional when use_record_time is used)>

# 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
Expand Down
14 changes: 13 additions & 1 deletion lib/fluent/plugin/in_rdkafka_group.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
66 changes: 66 additions & 0 deletions lib/fluent/plugin/kafka_plugin_util.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Comment thread
kenhys marked this conversation as resolved.
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
51 changes: 3 additions & 48 deletions lib/fluent/plugin/out_rdkafka2.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand All @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
Loading