From 951cd72a92329a3f150793dbdec54fd908a27f0f Mon Sep 17 00:00:00 2001 From: choi se Date: Sat, 23 Jul 2022 23:00:43 +0900 Subject: [PATCH] Support SASL Protocol of Kafka (#1092) * Support SASL Protocol of Kafka * Add comments * Fix Unused parameter detected in function RVV-B0012 --- input_kafka.go | 4 +-- kafka.go | 80 ++++++++++++++++++++++++++++++++++++++++++------- output_kafka.go | 4 +-- settings.go | 8 +++++ 4 files changed, 81 insertions(+), 15 deletions(-) diff --git a/input_kafka.go b/input_kafka.go index dc33e21..ebd10ed 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -19,8 +19,8 @@ type KafkaInput struct { } // NewKafkaInput creates instance of kafka consumer client with TLS config -func NewKafkaInput(address string, config *InputKafkaConfig, tlsConfig *KafkaTLSConfig) *KafkaInput { - c := NewKafkaConfig(tlsConfig) +func NewKafkaInput(_ string, config *InputKafkaConfig, tlsConfig *KafkaTLSConfig) *KafkaInput { + c := NewKafkaConfig(&config.SASLConfig, tlsConfig) var con sarama.Consumer diff --git a/kafka.go b/kafka.go index 46bc6c6..a73a812 100644 --- a/kafka.go +++ b/kafka.go @@ -2,6 +2,8 @@ package main import ( "bytes" + "crypto/sha256" + "crypto/sha512" "crypto/tls" "crypto/x509" "errors" @@ -11,25 +13,34 @@ import ( "github.com/Shopify/sarama" "github.com/buger/goreplay/proto" + "github.com/xdg-go/scram" ) +// SASLKafkaConfig SASL configuration +type SASLKafkaConfig struct { + UseSASL bool `json:"input-kafka-use-sasl"` + Mechanism string `json:"input-kafka-mechanism"` + Username string `json:"input-kafka-username"` + Password string `json:"input-kafka-password"` +} + // InputKafkaConfig should contains required information to // build producers. type InputKafkaConfig struct { - producer sarama.AsyncProducer - consumer sarama.Consumer - Host string `json:"input-kafka-host"` - Topic string `json:"input-kafka-topic"` - UseJSON bool `json:"input-kafka-json-format"` + consumer sarama.Consumer + Host string `json:"input-kafka-host"` + Topic string `json:"input-kafka-topic"` + UseJSON bool `json:"input-kafka-json-format"` + SASLConfig SASLKafkaConfig } // OutputKafkaConfig is the representation of kfka output configuration type OutputKafkaConfig struct { - producer sarama.AsyncProducer - consumer sarama.Consumer - Host string `json:"output-kafka-host"` - Topic string `json:"output-kafka-topic"` - UseJSON bool `json:"output-kafka-json-format"` + producer sarama.AsyncProducer + Host string `json:"output-kafka-host"` + Topic string `json:"output-kafka-topic"` + UseJSON bool `json:"output-kafka-json-format"` + SASLConfig SASLKafkaConfig } // KafkaTLSConfig should contains TLS certificates for connecting to secured Kafka clusters @@ -83,7 +94,7 @@ func NewTLSConfig(clientCertFile, clientKeyFile, caCertFile string) (*tls.Config } // NewKafkaConfig returns Kafka config with or without TLS -func NewKafkaConfig(tlsConfig *KafkaTLSConfig) *sarama.Config { +func NewKafkaConfig(saslConfig *SASLKafkaConfig, tlsConfig *KafkaTLSConfig) *sarama.Config { config := sarama.NewConfig() // Configuration options go here if tlsConfig != nil && (tlsConfig.ClientCert != "" || tlsConfig.CACert != "") { @@ -94,6 +105,18 @@ func NewKafkaConfig(tlsConfig *KafkaTLSConfig) *sarama.Config { } config.Net.TLS.Config = tlsConfig } + if saslConfig.UseSASL { + mechanism := sarama.SASLMechanism(saslConfig.Mechanism) + config.Net.SASL.Enable = saslConfig.UseSASL + config.Net.SASL.Mechanism = mechanism + config.Net.SASL.User = saslConfig.Username + config.Net.SASL.Password = saslConfig.Password + if mechanism == sarama.SASLTypeSCRAMSHA256 { + config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA256} } + } else if mechanism == sarama.SASLTypeSCRAMSHA512 { + config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA512} } + } + } return config } @@ -115,3 +138,38 @@ func (m KafkaMessage) Dump() ([]byte, error) { return b.Bytes(), nil } + +var ( + // SHA256 SASLMechanism + SHA256 scram.HashGeneratorFcn = sha256.New + // SHA512 SASLMechanism + SHA512 scram.HashGeneratorFcn = sha512.New +) + +// XDGSCRAMClient for SASL-Protocol +type XDGSCRAMClient struct { + *scram.Client + *scram.ClientConversation + scram.HashGeneratorFcn +} + +// Begin of XDGSCRAMClient +func (x *XDGSCRAMClient) Begin(userName, password, authzID string) (err error) { + x.Client, err = x.HashGeneratorFcn.NewClient(userName, password, authzID) + if err != nil { + return err + } + x.ClientConversation = x.Client.NewConversation() + return nil +} + +// Step of XDGSCRAMClient +func (x *XDGSCRAMClient) Step(challenge string) (response string, err error) { + response, err = x.ClientConversation.Step(challenge) + return +} + +// Done of XDGSCRAMClient +func (x *XDGSCRAMClient) Done() bool { + return x.ClientConversation.Done() +} diff --git a/output_kafka.go b/output_kafka.go index 816a7be..5398b79 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -23,8 +23,8 @@ type KafkaOutput struct { const KafkaOutputFrequency = 500 // NewKafkaOutput creates instance of kafka producer client with TLS config -func NewKafkaOutput(address string, config *OutputKafkaConfig, tlsConfig *KafkaTLSConfig) PluginWriter { - c := NewKafkaConfig(tlsConfig) +func NewKafkaOutput(_ string, config *OutputKafkaConfig, tlsConfig *KafkaTLSConfig) PluginWriter { + c := NewKafkaConfig(&config.SASLConfig, tlsConfig) var producer sarama.AsyncProducer diff --git a/settings.go b/settings.go index d49451b..e3d08ac 100644 --- a/settings.go +++ b/settings.go @@ -229,10 +229,18 @@ func init() { flag.StringVar(&Settings.OutputKafkaConfig.Host, "output-kafka-host", "", "Read request and response stats from Kafka:\n\tgor --input-raw :8080 --output-kafka-host '192.168.0.1:9092,192.168.0.2:9092'") flag.StringVar(&Settings.OutputKafkaConfig.Topic, "output-kafka-topic", "", "Read request and response stats from Kafka:\n\tgor --input-raw :8080 --output-kafka-topic 'kafka-log'") flag.BoolVar(&Settings.OutputKafkaConfig.UseJSON, "output-kafka-json-format", false, "If turned on, it will serialize messages from GoReplay text format to JSON.") + flag.BoolVar(&Settings.OutputKafkaConfig.SASLConfig.UseSASL, "output-kafka-use-sasl", false, "--output-kafka-use-sasl true") + flag.StringVar(&Settings.OutputKafkaConfig.SASLConfig.Mechanism, "output-kafka-mechanism", "", "mechanism\n\tgor --input-raw :8080 --output-kafka-mechanism 'SCRAM-SHA-512'") + flag.StringVar(&Settings.OutputKafkaConfig.SASLConfig.Username, "output-kafka-username", "", "username\n\tgor --input-raw :8080 --output-kafka-username 'username'") + flag.StringVar(&Settings.OutputKafkaConfig.SASLConfig.Password, "output-kafka-password", "", "password\n\tgor --input-raw :8080 --output-kafka-password 'password'") flag.StringVar(&Settings.InputKafkaConfig.Host, "input-kafka-host", "", "Send request and response stats to Kafka:\n\tgor --output-stdout --input-kafka-host '192.168.0.1:9092,192.168.0.2:9092'") flag.StringVar(&Settings.InputKafkaConfig.Topic, "input-kafka-topic", "", "Send request and response stats to Kafka:\n\tgor --output-stdout --input-kafka-topic 'kafka-log'") flag.BoolVar(&Settings.InputKafkaConfig.UseJSON, "input-kafka-json-format", false, "If turned on, it will assume that messages coming in JSON format rather than GoReplay text format.") + flag.BoolVar(&Settings.InputKafkaConfig.SASLConfig.UseSASL, "input-kafka-use-sasl", false, "use-sasl\n\t--use-sasl true") + flag.StringVar(&Settings.InputKafkaConfig.SASLConfig.Mechanism, "input-kafka-mechanism", "", "mechanism\n\tgor --input-raw :8080 --output-kafka-mechanism 'SCRAM-SHA-512'") + flag.StringVar(&Settings.InputKafkaConfig.SASLConfig.Username, "input-kafka-username", "", "username\n\tgor --input-raw :8080 --output-kafka-username 'username'") + flag.StringVar(&Settings.InputKafkaConfig.SASLConfig.Password, "input-kafka-password", "", "password\n\tgor --input-raw :8080 --output-kafka-password 'password'") flag.StringVar(&Settings.KafkaTLSConfig.CACert, "kafka-tls-ca-cert", "", "CA certificate for Kafka TLS Config:\n\tgor --input-raw :3000 --output-kafka-host '192.168.0.1:9092' --output-kafka-topic 'topic' --kafka-tls-ca-cert cacert.cer.pem --kafka-tls-client-cert client.cer.pem --kafka-tls-client-key client.key.pem") flag.StringVar(&Settings.KafkaTLSConfig.ClientCert, "kafka-tls-client-cert", "", "Client certificate for Kafka TLS Config (mandatory with to kafka-tls-ca-cert and kafka-tls-client-key)")