diff --git a/input_kafka.go b/input_kafka.go index 48c890f..609a704 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -3,9 +3,8 @@ package main import ( "encoding/json" "log" - "time" - "github.com/Shopify/sarama" + "github.com/Shopify/sarama/mocks" ) // KafkaInput is used for recieving Kafka messages and @@ -21,9 +20,17 @@ func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput { c := sarama.NewConfig() // Configuration options go here - con, err := sarama.NewConsumer([]string{config.host}, c) - if err != nil { - log.Fatalln("Failed to start Sarama(Kafka) consumer:", err) + var con sarama.Consumer + + if config.consumer.(*mocks.Consumer) != nil { + con = config.consumer + } else { + var err error + con, err = sarama.NewConsumer([]string{config.host}, c) + + if err != nil { + log.Fatalln("Failed to start Sarama(Kafka) consumer:", err) + } } partitions, err := con.Partitions(config.topic) @@ -72,21 +79,23 @@ func (i *KafkaInput) ErrorHandler(consumer sarama.PartitionConsumer) { func (i *KafkaInput) Read(data []byte) (int, error) { message := <-i.messages - var kafkaMessage KafkaMessage - json.Unmarshal(message.Value, &kafkaMessage) + if !i.config.useJSON { + copy(data, message.Value) + return len(message.Value), nil + } else { + var kafkaMessage KafkaMessage + json.Unmarshal(message.Value, &kafkaMessage) - buf, err := kafkaMessage.Dump() - if err != nil { - log.Println("Failed to decode access log entry:", err) - return 0, err + buf, err := kafkaMessage.Dump() + if err != nil { + log.Println("Failed to decode access log entry:", err) + return 0, err + } + + copy(data, buf) + + return len(buf), nil } - - header := payloadHeader(RequestPayload, uuid(), time.Now().UnixNano(), -1) - - copy(data[0:len(header)], header) - copy(data[len(header):], buf) - - return len(buf) + len(header), nil } func (i *KafkaInput) String() string { diff --git a/input_kafka_test.go b/input_kafka_test.go new file mode 100644 index 0000000..3a0803d --- /dev/null +++ b/input_kafka_test.go @@ -0,0 +1,61 @@ +package main + +import ( + "github.com/Shopify/sarama" + "github.com/Shopify/sarama/mocks" + "testing" +) + +func TestInputKafkaRAW(t *testing.T) { + consumer := mocks.NewConsumer(t, nil) + defer consumer.Close() + + consumer.ExpectConsumePartition("test", 0, mocks.AnyOffset).YieldMessage(&sarama.ConsumerMessage{Value: []byte("1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n")}) + consumer.SetTopicMetadata( + map[string][]int32{"test": []int32{0}}, + ) + + input := NewKafkaInput("", &KafkaConfig{ + consumer: consumer, + topic: "test", + useJSON: false, + }) + + buf := make([]byte, 1024) + n, err := input.Read(buf) + + if err != nil { + t.Fatal(err) + } + + if string(buf[:n]) != "1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n" { + t.Error("Message not properly decoded: ", string(buf[:n]), n) + } +} + +func TestInputKafkaJSON(t *testing.T) { + consumer := mocks.NewConsumer(t, nil) + defer consumer.Close() + + consumer.ExpectConsumePartition("test", 0, mocks.AnyOffset).YieldMessage(&sarama.ConsumerMessage{Value: []byte(`{"Req_URL":"/","Req_Type":"1","Req_ID":"2","Req_Ts":"3","Req_Method":"GET","Req_Headers":{"Header":"1"}}`)}) + consumer.SetTopicMetadata( + map[string][]int32{"test": []int32{0}}, + ) + + input := NewKafkaInput("", &KafkaConfig{ + consumer: consumer, + topic: "test", + useJSON: true, + }) + + buf := make([]byte, 1024) + n, err := input.Read(buf) + + if err != nil { + t.Fatal(err) + } + + if string(buf[:n]) != "1 2 3\nGET / HTTP/1.1\r\nHeader: 1\r\n\r\n" { + t.Error("Message not properly decoded: ", string(buf[:n]), n) + } +} diff --git a/kafka.go b/kafka.go index ef0849d..0c23032 100644 --- a/kafka.go +++ b/kafka.go @@ -3,21 +3,27 @@ package main import ( "bytes" "fmt" - + "github.com/Shopify/sarama" "github.com/buger/gor/proto" ) // KafkaConfig should contains required information to // build producers. type KafkaConfig struct { - host string - topic string + host string + topic string + producer sarama.AsyncProducer + consumer sarama.Consumer + useJSON bool } // KafkaMessage should contains catched request information that should be // passed as Json to Apache Kafka. type KafkaMessage struct { ReqURL string `json:"Req_URL"` + ReqType string `json:"Req_Type"` + ReqID string `json:"Req_ID"` + ReqTs string `json:"Req_Ts"` ReqMethod string `json:"Req_Method"` ReqBody string `json:"Req_Body,omitempty"` ReqHeaders map[string]string `json:"Req_Headers,omitempty"` @@ -28,6 +34,7 @@ type KafkaMessage struct { func (m KafkaMessage) Dump() ([]byte, error) { var b bytes.Buffer + b.WriteString(fmt.Sprintf("%s %s %s\n", m.ReqType, m.ReqID, m.ReqTs)) b.WriteString(fmt.Sprintf("%s %s HTTP/1.1", m.ReqMethod, m.ReqURL)) b.Write(proto.CLRF) for key, value := range m.ReqHeaders { diff --git a/output_kafka.go b/output_kafka.go index 35a0d67..336800f 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -2,13 +2,13 @@ package main import ( "encoding/json" + "github.com/Shopify/sarama" + "github.com/Shopify/sarama/mocks" + "github.com/buger/gor/proto" "io" "log" "strings" "time" - - "github.com/Shopify/sarama" - "github.com/buger/gor/proto" ) // KafkaOutput is used for sending payloads to kafka in JSON format. @@ -23,15 +23,23 @@ const KafkaOutputFrequency = 500 // NewKafkaOutput creates instance of kafka producer client. func NewKafkaOutput(address string, config *KafkaConfig) io.Writer { c := sarama.NewConfig() - c.Producer.RequiredAcks = sarama.WaitForLocal - c.Producer.Compression = sarama.CompressionSnappy - c.Producer.Flush.Frequency = KafkaOutputFrequency * time.Millisecond - brokerList := strings.Split(config.host, ",") + var producer sarama.AsyncProducer - producer, err := sarama.NewAsyncProducer(brokerList, c) - if err != nil { - log.Fatalln("Failed to start Sarama(Kafka) producer:", err) + if config.producer.(*mocks.AsyncProducer) != nil { + producer = config.producer + } else { + c.Producer.RequiredAcks = sarama.WaitForLocal + c.Producer.Compression = sarama.CompressionSnappy + c.Producer.Flush.Frequency = KafkaOutputFrequency * time.Millisecond + + brokerList := strings.Split(config.host, ",") + + var err error + producer, err = sarama.NewAsyncProducer(brokerList, c) + if err != nil { + log.Fatalln("Failed to start Sarama(Kafka) producer:", err) + } } o := &KafkaOutput{ @@ -55,22 +63,32 @@ func (o *KafkaOutput) ErrorHandler() { } func (o *KafkaOutput) Write(data []byte) (n int, err error) { - headers := make(map[string]string) - proto.ParseHeaders([][]byte{data}, func(header []byte, value []byte) bool { - headers[string(header)] = string(value) - return true - }) + var message sarama.StringEncoder - req := payloadBody(data) + if !o.config.useJSON { + message = sarama.StringEncoder(data) + } else { + headers := make(map[string]string) + proto.ParseHeaders([][]byte{data}, func(header []byte, value []byte) bool { + headers[string(header)] = string(value) + return true + }) - kafkaMessage := KafkaMessage{ - ReqURL: string(proto.Path(req)), - ReqMethod: string(proto.Method(req)), - ReqBody: string(proto.Body(req)), - ReqHeaders: headers, + meta := payloadMeta(data) + req := payloadBody(data) + + kafkaMessage := KafkaMessage{ + ReqURL: string(proto.Path(req)), + ReqType: string(meta[0]), + ReqID: string(meta[1]), + ReqTs: string(meta[2]), + ReqMethod: string(proto.Method(req)), + ReqBody: string(proto.Body(req)), + ReqHeaders: headers, + } + jsonMessage, _ := json.Marshal(&kafkaMessage) + message = sarama.StringEncoder(jsonMessage) } - jsonMessage, _ := json.Marshal(&kafkaMessage) - message := sarama.StringEncoder(jsonMessage) o.producer.Input() <- &sarama.ProducerMessage{ Topic: o.config.topic, diff --git a/output_kafka_test.go b/output_kafka_test.go new file mode 100644 index 0000000..cf9efe0 --- /dev/null +++ b/output_kafka_test.go @@ -0,0 +1,53 @@ +package main + +import ( + "github.com/Shopify/sarama" + "github.com/Shopify/sarama/mocks" + "testing" +) + +func TestOutputKafkaRAW(t *testing.T) { + config := sarama.NewConfig() + config.Producer.Return.Successes = true + producer := mocks.NewAsyncProducer(t, config) + producer.ExpectInputAndSucceed() + + output := NewKafkaOutput("", &KafkaConfig{ + producer: producer, + topic: "test", + useJSON: false, + }) + + output.Write([]byte("1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n")) + + resp := <-producer.Successes() + + data, _ := resp.Value.Encode() + + if string(data) != "1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n" { + t.Error("Message not properly encoded: ", string(data)) + } +} + +func TestOutputKafkaJSON(t *testing.T) { + config := sarama.NewConfig() + config.Producer.Return.Successes = true + producer := mocks.NewAsyncProducer(t, config) + producer.ExpectInputAndSucceed() + + output := NewKafkaOutput("", &KafkaConfig{ + producer: producer, + topic: "test", + useJSON: true, + }) + + output.Write([]byte("1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n")) + + resp := <-producer.Successes() + + data, _ := resp.Value.Encode() + + if string(data) != `{"Req_URL":"/","Req_Type":"1","Req_ID":"2","Req_Ts":"3","Req_Method":"GET","Req_Headers":{"Header":"1"}}` { + t.Error("Message not properly encoded: ", string(data)) + } +} diff --git a/settings.go b/settings.go index 90fa0af..62cdc79 100644 --- a/settings.go +++ b/settings.go @@ -131,9 +131,11 @@ 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.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.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")