diff --git a/capture/capture.go b/capture/capture.go index d029fb8..b06cbef 100644 --- a/capture/capture.go +++ b/capture/capture.go @@ -69,7 +69,7 @@ const ( // Set is here so that EngineType can implement flag.Var func (eng *EngineType) Set(v string) error { switch v { - case "", "libcap": + case "", "libpcap": *eng = EnginePcap case "pcap_file": *eng = EnginePcapFile diff --git a/input_kafka.go b/input_kafka.go index 6796c4e..80dd2db 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -19,8 +19,12 @@ type KafkaInput struct { // NewKafkaInput creates instance of kafka consumer client. func NewKafkaInput(address string, config *InputKafkaConfig) *KafkaInput { - c := sarama.NewConfig() - // Configuration options go here + return NewKafkaInputWithTLS(address, config, nil) +} + +// NewKafkaInputWithTLS creates instance of kafka consumer client with TLS +func NewKafkaInputWithTLS(address string, config *InputKafkaConfig, tlsConfig *KafkaTLSConfig) *KafkaInput { + c := NewKafkaConfig(tlsConfig) var con sarama.Consumer diff --git a/kafka.go b/kafka.go index f5163e1..29c55f1 100644 --- a/kafka.go +++ b/kafka.go @@ -2,7 +2,11 @@ package main import ( "bytes" + "crypto/tls" + "crypto/x509" "fmt" + "io/ioutil" + "log" "github.com/Shopify/sarama" "github.com/buger/goreplay/proto" @@ -27,6 +31,13 @@ type OutputKafkaConfig struct { UseJSON bool `json:"output-kafka-json-format"` } +// KafkaTLSConfig should contains TLS certificates for connecting to secured Kafka clusters +type KafkaTLSConfig struct { + CACert string `json:"kafka-tls-ca-cert"` + clientCert string `json:"kafka-tls-client-cert"` + clientKey string `json:"kafka-tls-client-key"` +} + // KafkaMessage should contains catched request information that should be // passed as Json to Apache Kafka. type KafkaMessage struct { @@ -39,6 +50,44 @@ type KafkaMessage struct { ReqHeaders map[string]string `json:"Req_Headers,omitempty"` } +// NewTLSConfig loads TLS certificates +func NewTLSConfig(clientCertFile, clientKeyFile, caCertFile string) (*tls.Config, error) { + tlsConfig := tls.Config{} + + // Load client cert + cert, err := tls.LoadX509KeyPair(clientCertFile, clientKeyFile) + if err != nil { + return &tlsConfig, err + } + tlsConfig.Certificates = []tls.Certificate{cert} + + // Load CA cert + caCert, err := ioutil.ReadFile(caCertFile) + if err != nil { + return &tlsConfig, err + } + caCertPool := x509.NewCertPool() + caCertPool.AppendCertsFromPEM(caCert) + tlsConfig.RootCAs = caCertPool + + return &tlsConfig, err +} + +// NewKafkaConfig returns Kafka config with or without TLS +func NewKafkaConfig(tlsConfig *KafkaTLSConfig) *sarama.Config { + config := sarama.NewConfig() + // Configuration options go here + if (tlsConfig != nil) && (tlsConfig.CACert != "") && (tlsConfig.clientCert != "") && (tlsConfig.clientKey != "") { + config.Net.TLS.Enable = true + tlsConfig, err := NewTLSConfig(tlsConfig.clientCert, tlsConfig.clientKey, tlsConfig.CACert) + if err != nil { + log.Fatal(err) + } + config.Net.TLS.Config = tlsConfig + } + return config +} + // Dump returns the given request in its HTTP/1.x wire // representation. func (m KafkaMessage) Dump() ([]byte, error) { diff --git a/output_kafka.go b/output_kafka.go index b94bc5e..d030ca2 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -24,7 +24,12 @@ const KafkaOutputFrequency = 500 // NewKafkaOutput creates instance of kafka producer client. func NewKafkaOutput(address string, config *OutputKafkaConfig) io.Writer { - c := sarama.NewConfig() + return NewKafkaOutputWithTLS(address, config, nil) +} + +// NewKafkaOutputWithTLS creates instance of kafka producer client. +func NewKafkaOutputWithTLS(address string, config *OutputKafkaConfig, tlsConfig *KafkaTLSConfig) io.Writer { + c := NewKafkaConfig(tlsConfig) var producer sarama.AsyncProducer diff --git a/plugins.go b/plugins.go index 032a862..151d189 100644 --- a/plugins.go +++ b/plugins.go @@ -137,11 +137,11 @@ func NewPlugins() *InOutPlugins { } if Settings.OutputKafkaConfig.Host != "" && Settings.OutputKafkaConfig.Topic != "" { - plugins.registerPlugin(NewKafkaOutput, "", &Settings.OutputKafkaConfig) + plugins.registerPlugin(NewKafkaOutput, "", &Settings.OutputKafkaConfig, &Settings.KafkaTLSConfig) } if Settings.InputKafkaConfig.Host != "" && Settings.InputKafkaConfig.Topic != "" { - plugins.registerPlugin(NewKafkaInput, "", &Settings.InputKafkaConfig) + plugins.registerPlugin(NewKafkaInput, "", &Settings.InputKafkaConfig, &Settings.KafkaTLSConfig) } return plugins diff --git a/settings.go b/settings.go index a5325b3..32d71aa 100644 --- a/settings.go +++ b/settings.go @@ -69,6 +69,7 @@ type AppSettings struct { InputKafkaConfig InputKafkaConfig OutputKafkaConfig OutputKafkaConfig + KafkaTLSConfig KafkaTLSConfig } // Settings holds Gor configuration @@ -177,14 +178,18 @@ func init() { flag.BoolVar(&Settings.OutputBinaryConfig.Debug, "output-binary-debug", false, "Enables binary debug output.") /* outputBinaryConfig */ - 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.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.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.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.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)") + flag.StringVar(&Settings.KafkaTLSConfig.clientKey, "kafka-tls-client-key", "", "Client Key for Kafka TLS Config (mandatory with to kafka-tls-client-cert and kafka-tls-client-key)") + flag.Var(&Settings.ModifierConfig.Headers, "http-set-header", "Inject additional headers to http reqest:\n\tgor --input-raw :8080 --output-http staging.com --http-set-header 'User-Agent: Gor'") flag.Var(&Settings.ModifierConfig.Headers, "output-http-header", "WARNING: `--output-http-header` DEPRECATED, use `--http-set-header` instead")