From 0831476beb90bb03e8d2648dfcc8d2f63f1486cb Mon Sep 17 00:00:00 2001 From: "Huu Khiem (Mark)" Date: Tue, 27 Sep 2016 15:40:26 +0800 Subject: [PATCH 01/20] Set default value for output-http-timeout at declaration (#371) The actual default is 5s, set in the call of NewHTTPClient, even if we use `--output-http--timeout=0`. It's clearer to user if this value id declared upfront. --- http_client.go | 4 ---- settings.go | 2 +- 2 files changed, 1 insertion(+), 5 deletions(-) diff --git a/http_client.go b/http_client.go index c4ec8f2..0d4df18 100644 --- a/http_client.go +++ b/http_client.go @@ -64,10 +64,6 @@ func NewHTTPClient(baseURL string, config *HTTPClientConfig) *HTTPClient { } } - if config.Timeout.Nanoseconds() == 0 { - config.Timeout = 5 * time.Second - } - config.ConnectionTimeout = config.Timeout if config.ResponseBufferSize == 0 { diff --git a/settings.go b/settings.go index 04828e7..b71bf56 100644 --- a/settings.go +++ b/settings.go @@ -118,7 +118,7 @@ func init() { flag.IntVar(&Settings.outputHTTPConfig.BufferSize, "output-http-response-buffer", 0, "HTTP response buffer size, all data after this size will be discarded.") flag.IntVar(&Settings.outputHTTPConfig.workers, "output-http-workers", 0, "Gor uses dynamic worker scaling by default. Enter a number to run a set number of workers.") flag.IntVar(&Settings.outputHTTPConfig.redirectLimit, "output-http-redirects", 0, "Enable how often redirects should be followed.") - flag.DurationVar(&Settings.outputHTTPConfig.Timeout, "output-http-timeout", 0, "Specify HTTP request/response timeout. By default 5s. Example: --output-http-timeout 30s") + flag.DurationVar(&Settings.outputHTTPConfig.Timeout, "output-http-timeout", 5 * time.Second, "Specify HTTP request/response timeout. By default 5s. Example: --output-http-timeout 30s") flag.BoolVar(&Settings.outputHTTPConfig.stats, "output-http-stats", false, "Report http output queue stats to console every 5 seconds.") flag.BoolVar(&Settings.outputHTTPConfig.OriginalHost, "http-original-host", false, "Normally gor replaces the Host http header with the host supplied with --output-http. This option disables that behavior, preserving the original Host header.") From 47db326b0f0e5dd30d3f432821a4cf8e27517211 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Tue, 11 Oct 2016 18:56:27 +0300 Subject: [PATCH 02/20] Fix truncated tcp check --- .gitignore | 3 +++ raw_socket_listener/listener.go | 2 +- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/.gitignore b/.gitignore index 4c1635a..6bc8d66 100644 --- a/.gitignore +++ b/.gitignore @@ -6,6 +6,7 @@ *.bin *.gz +*.zip *.class @@ -17,3 +18,5 @@ gor *.mprof *.pcap + +.DS_Store diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index 8e035e2..cecf979 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -465,7 +465,7 @@ func (t *Listener) readPcap() { } // Truncated TCP info - if len(data) < 13 { + if len(data) <= 13 { continue } From 22466d402818abdc935bfe0149b7569cee06df6e Mon Sep 17 00:00:00 2001 From: manjeshnilange Date: Tue, 25 Oct 2016 12:06:43 -0700 Subject: [PATCH 03/20] =?UTF-8?q?Adding=20ability=20to=20output=20http=20c?= =?UTF-8?q?lient=20to=20close=20connection=20on=20specific=20=E2=80=A6=20(?= =?UTF-8?q?#375)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Adding ability to output http client to close connection on specific response status * Addressed review comments --- http_client.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/http_client.go b/http_client.go index 0d4df18..b14242a 100644 --- a/http_client.go +++ b/http_client.go @@ -318,6 +318,11 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { } } + if bytes.Equal(proto.Status(payload), []byte("400")) { + c.Disconnect() + Debug("[HTTPClient] Closed connection on 400 response") + } + c.redirectsCount = 0 return payload, err From 76fb91966b77f722f5f155ffbcaa8bc6d1d1964b Mon Sep 17 00:00:00 2001 From: Jonathan Cremin Date: Thu, 27 Oct 2016 18:44:54 +0100 Subject: [PATCH 04/20] Treat PATCH as a HTTP verb --- proto/proto.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/proto/proto.go b/proto/proto.go index 0bdb819..53b627c 100644 --- a/proto/proto.go +++ b/proto/proto.go @@ -443,7 +443,7 @@ func Status(payload []byte) []byte { } var httpMethods []string = []string{ - "GET ", "OPTI", "HEAD", "POST", "PUT ", "DELE", "TRAC", "CONN" /* custom methods */, "BAN", "PURG", + "GET ", "OPTI", "HEAD", "POST", "PUT ", "DELE", "TRAC", "CONN", "PATC" /* custom methods */, "BAN", "PURG", } func IsHTTPPayload(payload []byte) bool { From e3a9a9b2c20b23448036fb566890e38c1e2fdd0b Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Mon, 31 Oct 2016 19:04:05 +0300 Subject: [PATCH 05/20] Add kafka output --- emitter.go | 2 +- output_kafka.go | 68 +++++++++++++++++++++++++++++++++++++++++++++++++ plugins.go | 2 ++ settings.go | 6 ++++- 4 files changed, 76 insertions(+), 2 deletions(-) create mode 100644 output_kafka.go diff --git a/emitter.go b/emitter.go index 1fc2a83..f085a2e 100644 --- a/emitter.go +++ b/emitter.go @@ -15,7 +15,7 @@ func Start(stop chan int) { middleware.ReadFrom(in) } - // We going only to read responses, so using same ReadFrom method + // We are going only to read responses, so using same ReadFrom method for _, out := range Plugins.Outputs { if r, ok := out.(io.Reader); ok { middleware.ReadFrom(r) diff --git a/output_kafka.go b/output_kafka.go new file mode 100644 index 0000000..cf2ea1f --- /dev/null +++ b/output_kafka.go @@ -0,0 +1,68 @@ +package main + +import ( + "github.com/Shopify/sarama" + "log" + "strings" + "time" +) + +// KafkaConfig should contains required information to +// build producers. +type KafkaConfig struct { + zookeeper string + topic string +} + +// KafkaOutput should make producer client. +type KafkaOutput struct { + address string + config *KafkaConfig + producer sarama.AsyncProducer +} + +// NewKafkaOutput creates instance of kafka producer client. +func NewKafkaOutput(address string, config *KafkaConfig) *KafkaOutput { + c := sarama.NewConfig() + c.Producer.RequiredAcks = sarama.WaitForLocal + c.Producer.Compression = sarama.CompressionSnappy + c.Producer.Flush.Frequency = 500 * time.Millisecond + + brokerList := strings.Split(config.zookeeper, ",") + + producer, err := sarama.NewAsyncProducer(brokerList, c) + if err != nil { + log.Fatalln("Failed to start Sarama(Kafka) producer:", err) + } + + o := &KafkaOutput{ + address: address, + config: config, + producer: producer, + } + + // Start infinite loop for tracking errors for kafka producer. + go o.ErrorHandler() + + return o +} + +// ErrorHandler should receive errors +func (o *KafkaOutput) ErrorHandler() { + for err := range o.producer.Errors() { + log.Println("Failed to write access log entry:", err) + } +} + +func (o *KafkaOutput) Write(data []byte) (n int, err error) { + buf := make(sarama.ByteEncoder, len(data)) + copy(buf, data) + + o.producer.Input() <- &sarama.ProducerMessage{ + Topic: o.config.topic, + Key: sarama.StringEncoder(o.address), + Value: buf, + } + + return len(data), nil +} diff --git a/plugins.go b/plugins.go index 1a0f6db..1d11971 100644 --- a/plugins.go +++ b/plugins.go @@ -142,4 +142,6 @@ func InitPlugins() { for _, options := range Settings.outputHTTP { registerPlugin(NewHTTPOutput, options, &Settings.outputHTTPConfig) } + + registerPlugin(NewKafkaOutput, &Settings.outputKafkaConfig) } diff --git a/settings.go b/settings.go index b71bf56..0b4973b 100644 --- a/settings.go +++ b/settings.go @@ -58,6 +58,8 @@ type AppSettings struct { outputHTTPConfig HTTPOutputConfig modifierConfig HTTPModifierConfig + + outputKafkaConfig KafkaConfig } // Settings holds Gor configuration @@ -118,13 +120,15 @@ func init() { flag.IntVar(&Settings.outputHTTPConfig.BufferSize, "output-http-response-buffer", 0, "HTTP response buffer size, all data after this size will be discarded.") flag.IntVar(&Settings.outputHTTPConfig.workers, "output-http-workers", 0, "Gor uses dynamic worker scaling by default. Enter a number to run a set number of workers.") flag.IntVar(&Settings.outputHTTPConfig.redirectLimit, "output-http-redirects", 0, "Enable how often redirects should be followed.") - flag.DurationVar(&Settings.outputHTTPConfig.Timeout, "output-http-timeout", 5 * time.Second, "Specify HTTP request/response timeout. By default 5s. Example: --output-http-timeout 30s") + flag.DurationVar(&Settings.outputHTTPConfig.Timeout, "output-http-timeout", 5*time.Second, "Specify HTTP request/response timeout. By default 5s. Example: --output-http-timeout 30s") flag.BoolVar(&Settings.outputHTTPConfig.stats, "output-http-stats", false, "Report http output queue stats to console every 5 seconds.") flag.BoolVar(&Settings.outputHTTPConfig.OriginalHost, "http-original-host", false, "Normally gor replaces the Host http header with the host supplied with --output-http. This option disables that behavior, preserving the original Host header.") flag.BoolVar(&Settings.outputHTTPConfig.Debug, "output-http-debug", false, "Enables http debug output.") flag.StringVar(&Settings.outputHTTPConfig.elasticSearch, "output-http-elasticsearch", "", "Send request and response stats to ElasticSearch:\n\tgor --input-raw :8080 --output-http staging.com --output-http-elasticsearch 'es_host:api_port/index_name'") + flag.StringVar(&Settings.outputKafkaConfig.zookeeper, "output-kafka-zookeeper", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-zookeeper '192.168.0.1:2181,192.168.0.2:2181'") + flag.StringVar(&Settings.outputKafkaConfig.topic, "output-kafka-topic", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-topic 'kafka-log'") 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") From 9b96dc20df2dd885c1653c00d5252b1ff25c6c3d Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Mon, 31 Oct 2016 19:29:59 +0300 Subject: [PATCH 06/20] Change Dockerfile because of issue on build echo.java --- Dockerfile | 4 +++- examples/middleware/echo.java | 24 ++++++++++++------------ 2 files changed, 15 insertions(+), 13 deletions(-) diff --git a/Dockerfile b/Dockerfile index 98fb7a4..f3d1874 100644 --- a/Dockerfile +++ b/Dockerfile @@ -18,5 +18,7 @@ RUN go get -u github.com/golang/lint/golint WORKDIR /go/src/github.com/buger/gor/ ADD . /go/src/github.com/buger/gor/ -RUN javac -cp /tmp/commons-io-2.4/commons-io-2.4.jar ./examples/middleware/echo.java +RUN wget http://archive.apache.org/dist/commons/io/binaries/commons-io-2.4-bin.tar.gz && tar xzf commons-io-2.4-bin.tar.gz && cd commons-io-2.4 && mv commons-io-2.4.jar /tmp/ +RUN wget http://archive.apache.org/dist/commons/codec/binaries/commons-codec-1.9-bin.tar.gz && tar xzf commons-codec-1.9-bin.tar.gz +RUN javac -cp commons-io-2.4/commons-io-2.4.jar -cp commons-codec-1.9/commons-codec-1.9.jar ./examples/middleware/echo.java RUN go get \ No newline at end of file diff --git a/examples/middleware/echo.java b/examples/middleware/echo.java index ffa885c..9fba333 100644 --- a/examples/middleware/echo.java +++ b/examples/middleware/echo.java @@ -6,19 +6,19 @@ import org.apache.commons.codec.DecoderException; import org.apache.commons.codec.binary.Hex; -public class Echo { - public static String decodeHexString(String s) throws DecoderException { - return new String(Hex.decodeHex(s.toCharArray())); - } +class Echo { + public static String decodeHexString(String s) throws DecoderException { + return new String(Hex.decodeHex(s.toCharArray())); + } - public static String encodeHexString(String s) { - return new String(Hex.encodeHex(s.getBytes())); - } + public static String encodeHexString(String s) { + return new String(Hex.encodeHex(s.getBytes())); + } - public static String transformHTTPMessage(String req) { - // do actual transformations here - return req; - } + public static String transformHTTPMessage(String req) { + // do actual transformations here + return req; + } public static void main(String[] args) throws DecoderException { if(args != null){ @@ -29,7 +29,7 @@ public class Echo { } BufferedReader stdin = new BufferedReader(new InputStreamReader( - System.in)); + System.in)); String line = null; try { From 41d3ce48a3d8e372935f09cfcfdf39f8d2259bb3 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Tue, 1 Nov 2016 17:51:15 +0300 Subject: [PATCH 07/20] Add json message output for kafka --- elasticsearch.go | 8 ++++--- examples/server.go | 15 ++++++++++++ output_kafka.go | 58 +++++++++++++++++++++++++++++++++++----------- plugins.go | 2 +- settings.go | 3 ++- 5 files changed, 67 insertions(+), 19 deletions(-) create mode 100644 examples/server.go diff --git a/elasticsearch.go b/elasticsearch.go index a626853..2635b89 100644 --- a/elasticsearch.go +++ b/elasticsearch.go @@ -84,9 +84,11 @@ func (p *ESPlugin) Init(URI string) { p.done = make(chan bool) p.indexor.Start() - // Only start the ErrorHandler goroutine when in verbose mode - // no need to burn ressources otherwise - go p.ErrorHandler() + if Settings.verbose { + // Only start the ErrorHandler goroutine when in verbose mode + // no need to burn ressources otherwise + go p.ErrorHandler() + } log.Println("Initialized Elasticsearch Plugin") return diff --git a/examples/server.go b/examples/server.go new file mode 100644 index 0000000..8be30c6 --- /dev/null +++ b/examples/server.go @@ -0,0 +1,15 @@ +package main + +import ( + "io" + "net/http" +) + +func hello(w http.ResponseWriter, r *http.Request) { + io.WriteString(w, "Hello world!") +} + +func main() { + http.HandleFunc("/", hello) + http.ListenAndServe(":8000", nil) +} diff --git a/output_kafka.go b/output_kafka.go index cf2ea1f..9af4ca0 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -1,7 +1,10 @@ package main import ( + "encoding/json" "github.com/Shopify/sarama" + "github.com/buger/gor/proto" + "io" "log" "strings" "time" @@ -10,25 +13,41 @@ import ( // KafkaConfig should contains required information to // build producers. type KafkaConfig struct { - zookeeper string - topic string + host string + topic string } // KafkaOutput should make producer client. type KafkaOutput struct { - address string config *KafkaConfig producer sarama.AsyncProducer } +// KafkaMessage should contains catched request information that should be +// passed as Json to Apache Kafka. +type KafkaMessage struct { + ReqURL string `json:"Req_URL"` + ReqMethod string `json:"Req_Method"` + ReqUserAgent string `json:"Req_User-Agent"` + ReqAcceptLanguage string `json:"Req_Accept-Language,omitempty"` + ReqAccept string `json:"Req_Accept,omitempty"` + ReqAcceptEncoding string `json:"Req_Accept-Encoding,omitempty"` + ReqIfModifiedSince string `json:"Req_If-Modified-Since,omitempty"` + ReqConnection string `json:"Req_Connection,omitempty"` + ReqCookies string `json:"Req_Cookies,omitempty"` +} + +// KafkaOutputFrequency in milliseconds +const KafkaOutputFrequency = 500 + // NewKafkaOutput creates instance of kafka producer client. -func NewKafkaOutput(address string, config *KafkaConfig) *KafkaOutput { +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 = 500 * time.Millisecond + c.Producer.Flush.Frequency = KafkaOutputFrequency * time.Millisecond - brokerList := strings.Split(config.zookeeper, ",") + brokerList := strings.Split(config.host, ",") producer, err := sarama.NewAsyncProducer(brokerList, c) if err != nil { @@ -36,13 +55,14 @@ func NewKafkaOutput(address string, config *KafkaConfig) *KafkaOutput { } o := &KafkaOutput{ - address: address, config: config, producer: producer, } - // Start infinite loop for tracking errors for kafka producer. - go o.ErrorHandler() + if Settings.verbose { + // Start infinite loop for tracking errors for kafka producer. + go o.ErrorHandler() + } return o } @@ -55,14 +75,24 @@ func (o *KafkaOutput) ErrorHandler() { } func (o *KafkaOutput) Write(data []byte) (n int, err error) { - buf := make(sarama.ByteEncoder, len(data)) - copy(buf, data) + kafkaMessage := KafkaMessage{ + ReqURL: string(proto.Path(data)), + ReqMethod: string(proto.Method(data)), + ReqUserAgent: string(proto.Header(data, []byte("User-Agent"))), + ReqAcceptLanguage: string(proto.Header(data, []byte("Accept-Language"))), + ReqAccept: string(proto.Header(data, []byte("Accept"))), + ReqAcceptEncoding: string(proto.Header(data, []byte("Accept-Encoding"))), + ReqIfModifiedSince: string(proto.Header(data, []byte("If-Modified-Since"))), + ReqConnection: string(proto.Header(data, []byte("Connection"))), + ReqCookies: string(proto.Header(data, []byte("Cookie"))), + } + jsonMessage, _ := json.Marshal(&kafkaMessage) + message := sarama.StringEncoder(jsonMessage) o.producer.Input() <- &sarama.ProducerMessage{ Topic: o.config.topic, - Key: sarama.StringEncoder(o.address), - Value: buf, + Value: message, } - return len(data), nil + return len(message), nil } diff --git a/plugins.go b/plugins.go index 1d11971..a625f9b 100644 --- a/plugins.go +++ b/plugins.go @@ -143,5 +143,5 @@ func InitPlugins() { registerPlugin(NewHTTPOutput, options, &Settings.outputHTTPConfig) } - registerPlugin(NewKafkaOutput, &Settings.outputKafkaConfig) + registerPlugin(NewKafkaOutput, "", &Settings.outputKafkaConfig) } diff --git a/settings.go b/settings.go index 0b4973b..ff417ca 100644 --- a/settings.go +++ b/settings.go @@ -127,7 +127,8 @@ func init() { flag.BoolVar(&Settings.outputHTTPConfig.Debug, "output-http-debug", false, "Enables http debug output.") flag.StringVar(&Settings.outputHTTPConfig.elasticSearch, "output-http-elasticsearch", "", "Send request and response stats to ElasticSearch:\n\tgor --input-raw :8080 --output-http staging.com --output-http-elasticsearch 'es_host:api_port/index_name'") - flag.StringVar(&Settings.outputKafkaConfig.zookeeper, "output-kafka-zookeeper", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-zookeeper '192.168.0.1:2181,192.168.0.2:2181'") + + flag.StringVar(&Settings.outputKafkaConfig.host, "output-kafka-host", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-host '192.168.0.1:2181,192.168.0.2:2181'") flag.StringVar(&Settings.outputKafkaConfig.topic, "output-kafka-topic", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-topic 'kafka-log'") 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'") From d424167f0496242ec7246ba6b6da3e1d8f9b9c64 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Tue, 1 Nov 2016 17:55:04 +0300 Subject: [PATCH 08/20] Change example of port for kafka host --- settings.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/settings.go b/settings.go index ff417ca..ac62f7c 100644 --- a/settings.go +++ b/settings.go @@ -128,7 +128,7 @@ func init() { flag.StringVar(&Settings.outputHTTPConfig.elasticSearch, "output-http-elasticsearch", "", "Send request and response stats to ElasticSearch:\n\tgor --input-raw :8080 --output-http staging.com --output-http-elasticsearch 'es_host:api_port/index_name'") - flag.StringVar(&Settings.outputKafkaConfig.host, "output-kafka-host", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-host '192.168.0.1:2181,192.168.0.2:2181'") + flag.StringVar(&Settings.outputKafkaConfig.host, "output-kafka-host", "", "Send request and response stats to 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", "", "Send request and response stats to Kafka:\n\tgor --input-raw :8080 --output-kafka-topic 'kafka-log'") 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'") From 0a038575f3d8eeca3158979a4303d4fc4fbeaa4f Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Wed, 2 Nov 2016 00:03:12 +0300 Subject: [PATCH 09/20] Add all headers and properly passed body --- output_kafka.go | 34 ++++++++++++++++------------------ 1 file changed, 16 insertions(+), 18 deletions(-) diff --git a/output_kafka.go b/output_kafka.go index 9af4ca0..4ebfc4c 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -26,15 +26,10 @@ type KafkaOutput struct { // KafkaMessage should contains catched request information that should be // passed as Json to Apache Kafka. type KafkaMessage struct { - ReqURL string `json:"Req_URL"` - ReqMethod string `json:"Req_Method"` - ReqUserAgent string `json:"Req_User-Agent"` - ReqAcceptLanguage string `json:"Req_Accept-Language,omitempty"` - ReqAccept string `json:"Req_Accept,omitempty"` - ReqAcceptEncoding string `json:"Req_Accept-Encoding,omitempty"` - ReqIfModifiedSince string `json:"Req_If-Modified-Since,omitempty"` - ReqConnection string `json:"Req_Connection,omitempty"` - ReqCookies string `json:"Req_Cookies,omitempty"` + ReqURL string `json:"Req_URL"` + ReqMethod string `json:"Req_Method"` + ReqBody string `json:"Req_Body,omitempty"` + ReqHeaders map[string]string `json:"Req_Headers,omitempty"` } // KafkaOutputFrequency in milliseconds @@ -75,16 +70,19 @@ 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 + }) + + req := payloadBody(data) + kafkaMessage := KafkaMessage{ - ReqURL: string(proto.Path(data)), - ReqMethod: string(proto.Method(data)), - ReqUserAgent: string(proto.Header(data, []byte("User-Agent"))), - ReqAcceptLanguage: string(proto.Header(data, []byte("Accept-Language"))), - ReqAccept: string(proto.Header(data, []byte("Accept"))), - ReqAcceptEncoding: string(proto.Header(data, []byte("Accept-Encoding"))), - ReqIfModifiedSince: string(proto.Header(data, []byte("If-Modified-Since"))), - ReqConnection: string(proto.Header(data, []byte("Connection"))), - ReqCookies: string(proto.Header(data, []byte("Cookie"))), + ReqURL: string(proto.Path(req)), + ReqMethod: string(proto.Method(req)), + ReqBody: string(proto.Body(req)), + ReqHeaders: headers, } jsonMessage, _ := json.Marshal(&kafkaMessage) message := sarama.StringEncoder(jsonMessage) From 40f7facc3328b6c9ba4dd9ab092119ed437cf937 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Wed, 2 Nov 2016 15:34:42 +0300 Subject: [PATCH 10/20] User-Agent could contains ':' inside of value. --- proto/proto.go | 7 ++++++- proto/proto_test.go | 22 ++++++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/proto/proto.go b/proto/proto.go index 0bdb819..c0bdd50 100644 --- a/proto/proto.go +++ b/proto/proto.go @@ -190,6 +190,7 @@ func ParseHeaders(payloads [][]byte, cb func(header []byte, value []byte) bool) i := 0 pIdx := 0 lineBreaks := 0 + newLineBreak := true for { if len(payloads)-1 < pIdx { @@ -206,6 +207,7 @@ func ParseHeaders(payloads [][]byte, cb func(header []byte, value []byte) bool) switch p[i] { case '\r', '\n': + newLineBreak = true lineBreaks++ // End of headers @@ -254,7 +256,10 @@ func ParseHeaders(payloads [][]byte, cb func(header []byte, value []byte) bool) hS = [2]int{-1, -1} hE = [2]int{-1, -1} case ':': - hE = [2]int{pIdx, i} + if newLineBreak { + hE = [2]int{pIdx, i} + newLineBreak = false + } default: lineBreaks = 0 diff --git a/proto/proto_test.go b/proto/proto_test.go index 32da8c9..7d25a38 100644 --- a/proto/proto_test.go +++ b/proto/proto_test.go @@ -139,11 +139,33 @@ func TestParseHeaders(t *testing.T) { "Host": "www.w3.org", "User-Agent": "Chrome", } + if !reflect.DeepEqual(headers, expected) { t.Error("Headers do not properly parsed", headers) } } +func TestParseHeadersWithComplexUserAgent(t *testing.T) { + // User-Agent could contain inside ':' + // Parser should wait for \r\n + payload := [][]byte{[]byte("POST /post HTTP/1.1\r\nContent-Length: 7\r\nHost: www.w3.or"), []byte("g\r\nUser-Ag"), []byte("ent:Mozilla/5.0 (Windows NT 6.1; WOW64; Trident/7.0; rv:11.0) like Gecko\r\n\r\n"), []byte("Fake-Header: asda")} + + headers := make(map[string]string) + + ParseHeaders(payload, func(header []byte, value []byte) bool { + headers[string(header)] = string(value) + return true + }) + + expected := map[string]string{ + "User-Agent": "Mozilla/5.0 (Windows NT 6.1; WOW64; Trident/7.0; rv:11.0) like Gecko", + } + + if expected["User-Agent"] != headers["User-Agent"] { + t.Errorf("Header 'User-Agent' expected '%s' and parsed: '%s'", expected["User-Agent"], headers["User-Agent"]) + } +} + func TestHeaderEquals(t *testing.T) { tests := []struct { h1 string From 638ddfb144b275c3adde9096229c4e6cdcd3d226 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Wed, 2 Nov 2016 17:50:45 +0300 Subject: [PATCH 11/20] Add one more test case to cover ':' inside of header value --- proto/proto_test.go | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/proto/proto_test.go b/proto/proto_test.go index 7d25a38..800799b 100644 --- a/proto/proto_test.go +++ b/proto/proto_test.go @@ -166,6 +166,37 @@ func TestParseHeadersWithComplexUserAgent(t *testing.T) { } } +func TestParseHeadersWithOrigin(t *testing.T) { + // User-Agent could contain inside ':' + // Parser should wait for \r\n + payload := [][]byte{[]byte("POST /post HTTP/1.1\r\nContent-Length: 7\r\nHost: www.w3.or"), []byte("g\r\nReferrer: http://127.0.0.1:3000\r\nOrigi"), []byte("n: https://www.example.com\r\nUser-Ag"), []byte("ent:Mozilla/5.0 (Windows NT 6.1; WOW64; Trident/7.0; rv:11.0) like Gecko\r\n\r\n"), []byte("in:https://www.example.com\r\n\r\n"), []byte("Fake-Header: asda")} + + headers := make(map[string]string) + + ParseHeaders(payload, func(header []byte, value []byte) bool { + headers[string(header)] = string(value) + return true + }) + + expected := map[string]string{ + "Origin": "https://www.example.com", + "User-Agent": "Mozilla/5.0 (Windows NT 6.1; WOW64; Trident/7.0; rv:11.0) like Gecko", + "Referrer": "http://127.0.0.1:3000", + } + + if expected["Referrer"] != headers["Referrer"] { + t.Errorf("Header 'Referrer' expected '%s' and parsed: '%s'", expected["Referrer"], headers["Referrer"]) + } + + if expected["Origin"] != headers["Origin"] { + t.Errorf("Header 'Origin' expected '%s' and parsed: '%s'", expected["Origin"], headers["Origin"]) + } + + if expected["User-Agent"] != headers["User-Agent"] { + t.Errorf("Header 'User-Agent' expected '%s' and parsed: '%s'", expected["User-Agent"], headers["User-Agent"]) + } +} + func TestHeaderEquals(t *testing.T) { tests := []struct { h1 string From f90dcc763b598bf4060af0fe05f13d9caafa4d57 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Tue, 8 Nov 2016 22:21:35 +0300 Subject: [PATCH 12/20] Remove server.go because of having 'gor file-server :8080' --- examples/server.go | 15 --------------- 1 file changed, 15 deletions(-) delete mode 100644 examples/server.go diff --git a/examples/server.go b/examples/server.go deleted file mode 100644 index 8be30c6..0000000 --- a/examples/server.go +++ /dev/null @@ -1,15 +0,0 @@ -package main - -import ( - "io" - "net/http" -) - -func hello(w http.ResponseWriter, r *http.Request) { - io.WriteString(w, "Hello world!") -} - -func main() { - http.HandleFunc("/", hello) - http.ListenAndServe(":8000", nil) -} From a1174a159c42361b0b151280631598d51548edf0 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Mon, 14 Nov 2016 17:47:30 +0300 Subject: [PATCH 13/20] Do not add port to Host header for https --- http_client.go | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/http_client.go b/http_client.go index b14242a..a47e850 100644 --- a/http_client.go +++ b/http_client.go @@ -58,11 +58,6 @@ func NewHTTPClient(baseURL string, config *HTTPClientConfig) *HTTPClient { } u, _ := url.Parse(baseURL) - if !strings.Contains(u.Host, ":") { - if u.Scheme != "http" { - u.Host += ":" + defaultPorts[u.Scheme] - } - } config.ConnectionTimeout = config.Timeout @@ -88,7 +83,7 @@ func (c *HTTPClient) Connect() (err error) { c.Disconnect() if !strings.Contains(c.host, ":") { - c.conn, err = net.DialTimeout("tcp", c.host+":80", c.config.ConnectionTimeout) + c.conn, err = net.DialTimeout("tcp", c.host + ":" + defaultPorts[c.scheme], c.config.ConnectionTimeout) } else { c.conn, err = net.DialTimeout("tcp", c.host, c.config.ConnectionTimeout) } From 963f7a73d7697d946314b7dd0d58310dcfaa6395 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Solano=20G=C3=B3mez?= Date: Wed, 23 Nov 2016 09:55:54 -0500 Subject: [PATCH 14/20] Add basic SNI support By adding the server name as part of the TLS client configuration, gor will now support connecting to hosts that require SNI (such as Amazon API Gateway). --- http_client.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/http_client.go b/http_client.go index a47e850..8ee42e4 100644 --- a/http_client.go +++ b/http_client.go @@ -89,7 +89,7 @@ func (c *HTTPClient) Connect() (err error) { } if c.scheme == "https" { - tlsConn := tls.Client(c.conn, &tls.Config{InsecureSkipVerify: true}) + tlsConn := tls.Client(c.conn, &tls.Config{InsecureSkipVerify: true, ServerName: c.host}) if err = tlsConn.Handshake(); err != nil { return From 12043e5a017ab2e9c4e2312c93f4466e9fbe8fb3 Mon Sep 17 00:00:00 2001 From: Yohan Legat Date: Thu, 1 Dec 2016 12:07:00 +0100 Subject: [PATCH 15/20] Resolve buger/gor#394 : Timeout for NO_CONTENT responses 1xx, 204 and 304 HTTP responses MUST NOT include a message body (see [RFC-2616](https://tools.ietf.org/html/rfc2616#section-4.4)). Also, a server MAY send a Content-Length header field in a 304 (Not Modified) and MUST NOT send a Content-Length header field in any response with a status code of 1xx (Informational) or 204 (No Content) (see [RFC-7230](https://tools.ietf.org/html/rfc7230#section-3.3.2)) --- http_client.go | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/http_client.go b/http_client.go index 8ee42e4..8007055 100644 --- a/http_client.go +++ b/http_client.go @@ -207,9 +207,14 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { if bytes.Equal(proto.Header(c.respBuf, []byte("Transfer-Encoding")), []byte("chunked")) { chunked = true } else { - l := proto.Header(c.respBuf, []byte("Content-Length")) - if len(l) > 0 { - contentLength, _ = strconv.Atoi(string(l)) + status, _ := strconv.Atoi(string(proto.Status(c.respBuf))) + if (status >= 100 && status < 200) || status == 204 || status == 304 { + contentLength = 0 + } else { + l := proto.Header(c.respBuf, []byte("Content-Length")) + if len(l) > 0 { + contentLength, _ = strconv.Atoi(string(l)) + } } } From bea255ab113f631d3bf1d33add36dc7a0e2ef046 Mon Sep 17 00:00:00 2001 From: Alexandr Korsak Date: Thu, 1 Dec 2016 16:46:18 +0300 Subject: [PATCH 16/20] Register kafka output plugin in case of having passed settings --- plugins.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/plugins.go b/plugins.go index a625f9b..6ee1516 100644 --- a/plugins.go +++ b/plugins.go @@ -143,5 +143,7 @@ func InitPlugins() { registerPlugin(NewHTTPOutput, options, &Settings.outputHTTPConfig) } - registerPlugin(NewKafkaOutput, "", &Settings.outputKafkaConfig) + if Settings.outputKafkaConfig.host != "" && Settings.outputKafkaConfig.topic != "" { + registerPlugin(NewKafkaOutput, "", &Settings.outputKafkaConfig) + } } From fa0e96b7efe0c345f230286ac09b2aaf558ee64e Mon Sep 17 00:00:00 2001 From: Yohan Legat Date: Tue, 29 Nov 2016 14:07:41 +0100 Subject: [PATCH 17/20] Resolve buger/gor#392 : use pcap timestamp Timestamp requests are now fetched from pcap payload --- raw_socket_listener/listener.go | 54 +++++++++++------- raw_socket_listener/listener_test.go | 76 ++++++++++++------------- raw_socket_listener/tcp_message.go | 9 ++- raw_socket_listener/tcp_message_test.go | 57 ++++++++++++------- raw_socket_listener/tcp_packet.go | 35 +++++++----- 5 files changed, 135 insertions(+), 96 deletions(-) diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index cecf979..1b1c8c0 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -33,6 +33,12 @@ import ( var _ = fmt.Println +type Packet struct { + srcIP []byte + data []byte + timestamp time.Time +} + // Listener handle traffic capture type Listener struct { mu sync.Mutex @@ -53,7 +59,7 @@ type Listener struct { respWithoutReq map[uint32]tcpID // Messages ready to be send to client - packetsChan chan []byte + packetsChan chan *Packet // Messages ready to be send to client messagesChan chan *TCPMessage @@ -88,7 +94,7 @@ const ( func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration) (l *Listener) { l = &Listener{} - l.packetsChan = make(chan []byte, 10000) + l.packetsChan = make(chan *Packet, 10000) l.messagesChan = make(chan *TCPMessage, 10000) l.quit = make(chan bool) l.readyCh = make(chan bool, 1) @@ -137,9 +143,9 @@ func (t *Listener) listen() { t.conn.Close() } return - case data := <-t.packetsChan: - packet := ParseTCPPacket(data[:16], data[16:]) - t.processTCPPacket(packet) + case packet := <-t.packetsChan: + tcpPacket := ParseTCPPacket(packet.srcIP, packet.data, packet.timestamp) + t.processTCPPacket(tcpPacket) case <-gcTicker: now := time.Now() @@ -522,11 +528,13 @@ func (t *Listener) readPcap() { } } - newBuf := make([]byte, len(data)+16) - copy(newBuf[:16], srcIP) - copy(newBuf[16:], data) + packetSrcIP := make([]byte, 16) + packetData := make([]byte, len(data)) - t.packetsChan <- newBuf + copy(packetSrcIP, srcIP) + copy(packetData, data) + + t.packetsChan <- t.buildPacket(srcIP, data, packet.Metadata().Timestamp) } } }(d) @@ -589,11 +597,7 @@ func (t *Listener) readPcapFile() { continue } - newBuf := make([]byte, len(data)+16) - copy(newBuf[:16], addr) - copy(newBuf[16:], data) - - t.packetsChan <- newBuf + t.packetsChan <- t.buildPacket(addr, data, packet.Metadata().Timestamp) } } } @@ -626,16 +630,26 @@ func (t *Listener) readRAWSocket() { if n > 0 { if t.isValidPacket(buf[:n]) { - newBuf := make([]byte, n+16) - copy(newBuf[16:], buf[:n]) - copy(newBuf[:16], []byte(addr.(*net.IPAddr).IP)) - - t.packetsChan <- newBuf + t.packetsChan <- t.buildPacket([]byte(addr.(*net.IPAddr).IP), buf[:n], time.Now()) } } } } +func (t *Listener) buildPacket(packetSrcIP []byte, packetData []byte, timestamp time.Time) *Packet { + copyPacketSrcIP := make([]byte, 16) + copyPacketData := make([]byte, len(packetData)) + + copy(copyPacketSrcIP, packetSrcIP) + copy(copyPacketData, packetSrcIP) + + return &Packet { + srcIP: packetSrcIP, + data: packetData, + timestamp:timestamp, + } +} + func (t *Listener) isValidPacket(buf []byte) bool { // To avoid full packet parsing every time, we manually parsing values needed for packet filtering // http://en.wikipedia.org/wiki/Transmission_Control_Protocol @@ -718,7 +732,7 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { message, ok := t.messages[packet.ID] if !ok { - message = NewTCPMessage(packet.Seq, packet.Ack, isIncoming) + message = NewTCPMessage(packet.Seq, packet.Ack, isIncoming, packet.timestamp) t.messages[packet.ID] = message if !isIncoming { diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 72e01de..8ade26e 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -15,10 +15,10 @@ func TestRawListenerInput(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) respAck := reqPacket.Seq + uint32(len(reqPacket.Data)) - respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) listener.packetsChan <- reqPacket.Dump() listener.packetsChan <- respPacket.Dump() @@ -52,11 +52,11 @@ func TestRawListenerInputResponseByClose(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) respAck := reqPacket.Seq + uint32(len(reqPacket.Data)) - respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\nConnection: close\r\n\r\nasd")) - finPacket := buildPacket(false, respAck, reqPacket.Seq+2, []byte("")) + respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\nConnection: close\r\n\r\nasd"), time.Now()) + finPacket := buildPacket(false, respAck, reqPacket.Seq+2, []byte(""), time.Now()) finPacket.IsFIN = true listener.packetsChan <- reqPacket.Dump() @@ -92,7 +92,7 @@ func TestRawListenerInputWithoutResponse(t *testing.T) { listener := NewListener("", "0", EnginePcap, false, 10*time.Millisecond) defer listener.Close() - reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) listener.packetsChan <- reqPacket.Dump() @@ -114,8 +114,8 @@ func TestRawListenerResponse(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n")) - respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) + respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) // If response packet comes before request listener.packetsChan <- respPacket.Dump() @@ -152,15 +152,15 @@ func TestShort100Continue(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\nExpect: 100-continue\r\n\r\n")) + reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\nExpect: 100-continue\r\n\r\n"), time.Now()) // Packet with data have different Seq - reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a"), time.Now()) + reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b"), time.Now()) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n")) + respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n"), time.Now()) // panic(int(uint32(len(reqPacket1.Data)) + uint32(len(reqPacket2.Data)) + uint32(len(reqPacket3.Data)))) - respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") @@ -172,15 +172,15 @@ func Test100ContinueWrongOrder(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\nExpect: 100-continue\r\n\r\n")) + reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\nExpect: 100-continue\r\n\r\n"), time.Now()) // Packet with data have different Seq - reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a"), time.Now()) + reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b"), time.Now()) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n")) + respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n"), time.Now()) // panic(int(uint32(len(reqPacket1.Data)) + uint32(len(reqPacket2.Data)) + uint32(len(reqPacket3.Data)))) - respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") @@ -191,15 +191,15 @@ func TestAlt100ContinueHeaderOrder(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 2\r\n\r\n")) + reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 2\r\n\r\n"), time.Now()) // Packet with data have different Seq - reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("a"), time.Now()) + reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b"), time.Now()) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n")) + respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n"), time.Now()) // panic(int(uint32(len(reqPacket1.Data)) + uint32(len(reqPacket2.Data)) + uint32(len(reqPacket3.Data)))) - respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket2 := buildPacket(false, reqPacket3.Seq+1 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") @@ -352,16 +352,16 @@ func TestRawListenerChunkedWrongOrder(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond) defer listener.Close() - reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n")) + reqPacket1 := buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n"), time.Now()) // Packet with data have different Seq - reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("1\r\na\r\n")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+uint32(len(reqPacket2.Data)), []byte("1\r\nb\r\n")) - reqPacket4 := buildPacket(true, 2, reqPacket3.Seq+uint32(len(reqPacket3.Data)), []byte("0\r\n\r\n")) + reqPacket2 := buildPacket(true, 2, reqPacket1.Seq+uint32(len(reqPacket1.Data)), []byte("1\r\na\r\n"), time.Now()) + reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+uint32(len(reqPacket2.Data)), []byte("1\r\nb\r\n"), time.Now()) + reqPacket4 := buildPacket(true, 2, reqPacket3.Seq+uint32(len(reqPacket3.Data)), []byte("0\r\n\r\n"), time.Now()) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n")) + respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n"), time.Now()) // panic(int(uint32(len(reqPacket1.Data)) + uint32(len(reqPacket2.Data)) + uint32(len(reqPacket3.Data)))) - respPacket2 := buildPacket(false, reqPacket4.Seq+5 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket2 := buildPacket(false, reqPacket4.Seq+5 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) // Should re-construct message from all possible combinations for i := 0; i < 6*5*4*3*2*1; i++ { @@ -381,13 +381,13 @@ func chunkedPostMessage() []*TCPPacket { ack := uint32(rand.Int63()) seq := uint32(rand.Int63()) - reqPacket1 := buildPacket(true, ack, seq, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n")) + reqPacket1 := buildPacket(true, ack, seq, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n"), time.Now()) // Packet with data have different Seq - reqPacket2 := buildPacket(true, ack, seq+47, []byte("1\r\na\r\n")) - reqPacket3 := buildPacket(true, ack, reqPacket2.Seq+5, []byte("1\r\nb\r\n")) - reqPacket4 := buildPacket(true, ack, reqPacket3.Seq+5, []byte("0\r\n\r\n")) + reqPacket2 := buildPacket(true, ack, seq+47, []byte("1\r\na\r\n"), time.Now()) + reqPacket3 := buildPacket(true, ack, reqPacket2.Seq+5, []byte("1\r\nb\r\n"), time.Now()) + reqPacket4 := buildPacket(true, ack, reqPacket3.Seq+5, []byte("0\r\n\r\n"), time.Now()) - respPacket := buildPacket(false, reqPacket4.Seq+5 /* len of data */, ack, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket := buildPacket(false, reqPacket4.Seq+5 /* len of data */, ack, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) return []*TCPPacket{ reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket, @@ -409,8 +409,8 @@ func postMessage() []*TCPPacket { } return []*TCPPacket{ - buildPacket(true, ack, seq, data), - buildPacket(false, seq+uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n\r\n")), + buildPacket(true, ack, seq, data, time.Now()), + buildPacket(false, seq+uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()), } } @@ -420,8 +420,8 @@ func getMessage() []*TCPPacket { seq := uint32(rand.Int63()) return []*TCPPacket{ - buildPacket(true, ack, seq, []byte("GET / HTTP/1.1\r\n\r\n")), - buildPacket(false, seq+18, seq2, []byte("HTTP/1.1 200 OK\r\n\r\n")), + buildPacket(true, ack, seq, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()), + buildPacket(false, seq+18, seq2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()), } } diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 5cb566e..5319594 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -48,9 +48,8 @@ type TCPMessage struct { } // NewTCPMessage pointer created from a Acknowledgment number and a channel of messages readuy to be deleted -func NewTCPMessage(Seq, Ack uint32, IsIncoming bool) (msg *TCPMessage) { - msg = &TCPMessage{Seq: Seq, Ack: Ack, IsIncoming: IsIncoming} - msg.Start = time.Now() +func NewTCPMessage(Seq, Ack uint32, IsIncoming bool, timestamp time.Time) (msg *TCPMessage) { + msg = &TCPMessage{Seq: Seq, Ack: Ack, IsIncoming: IsIncoming, Start: timestamp} return } @@ -138,6 +137,10 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { if packet.OrigAck != 0 { t.DataAck = packet.OrigAck } + + if packet.timestamp.Before(t.Start) { + t.Start = packet.timestamp + } } t.checkSeqIntegrity() diff --git a/raw_socket_listener/tcp_message_test.go b/raw_socket_listener/tcp_message_test.go index 781a824..1ce0d19 100644 --- a/raw_socket_listener/tcp_message_test.go +++ b/raw_socket_listener/tcp_message_test.go @@ -5,9 +5,10 @@ import ( "encoding/binary" _ "log" "testing" + "time" ) -func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPacket) { +func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte, timestamp time.Time) (packet *TCPPacket) { var srcPort, destPort uint16 // For tests `listening` port is 0 @@ -25,7 +26,7 @@ func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPack buf[12] = 64 buf = append(buf, Data...) - packet = ParseTCPPacket([]byte("123"), buf) + packet = ParseTCPPacket([]byte("123"), buf, timestamp) return packet } @@ -36,31 +37,31 @@ func buildMessage(p *TCPPacket) *TCPMessage { isIncoming = true } - m := NewTCPMessage(p.Seq, p.Ack, isIncoming) + m := NewTCPMessage(p.Seq, p.Ack, isIncoming, p.timestamp) m.AddPacket(p) return m } func TestTCPMessagePacketsOrder(t *testing.T) { - msg := buildMessage(buildPacket(true, 1, 1, []byte("a"))) - msg.AddPacket(buildPacket(true, 1, 2, []byte("b"))) + msg := buildMessage(buildPacket(true, 1, 1, []byte("a"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 2, []byte("b"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("ab")) { t.Error("Should contatenate packets in right order") } // When first packet have wrong order (Seq) - msg = buildMessage(buildPacket(true, 1, 2, []byte("b"))) - msg.AddPacket(buildPacket(true, 1, 1, []byte("a"))) + msg = buildMessage(buildPacket(true, 1, 2, []byte("b"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("a"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("ab")) { t.Error("Should contatenate packets in right order") } // Should ignore packets with same sequence - msg = buildMessage(buildPacket(true, 1, 1, []byte("a"))) - msg.AddPacket(buildPacket(true, 1, 1, []byte("a"))) + msg = buildMessage(buildPacket(true, 1, 1, []byte("a"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("a"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("a")) { t.Error("Should ignore packet with same Seq") @@ -68,8 +69,8 @@ func TestTCPMessagePacketsOrder(t *testing.T) { } func TestTCPMessageSize(t *testing.T) { - msg := buildMessage(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"))) - msg.AddPacket(buildPacket(true, 1, 2, []byte("b"))) + msg := buildMessage(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 2, []byte("b"), time.Now())) if msg.BodySize() != 2 { t.Error("Should count only body", msg.BodySize()) @@ -110,7 +111,7 @@ func TestTCPMessageIsComplete(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if tc.assocMessage { msg.AssocMessage = &TCPMessage{} } @@ -123,9 +124,9 @@ func TestTCPMessageIsComplete(t *testing.T) { } func TestTCPMessageIsSeqMissing(t *testing.T) { - p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n")) - p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n")) - p3 := buildPacket(false, 1, p2.Seq+uint32(len(p2.Data)), []byte("a")) + p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n"), time.Now()) + p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n"), time.Now()) + p3 := buildPacket(false, 1, p2.Seq+uint32(len(p2.Data)), []byte("a"), time.Now()) msg := buildMessage(p1) if msg.seqMissing { @@ -144,8 +145,8 @@ func TestTCPMessageIsSeqMissing(t *testing.T) { } func TestTCPMessageIsHeadersReceived(t *testing.T) { - p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n\r\n")) - p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n")) + p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) + p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n"), time.Now()) msg := buildMessage(p1) if msg.headerPacket == -1 { @@ -157,7 +158,7 @@ func TestTCPMessageIsHeadersReceived(t *testing.T) { t.Error("Should found double new line: headers received") } - msg = buildMessage(buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\nContent-Length: 1\r\n"))) + msg = buildMessage(buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\nContent-Length: 1\r\n"), time.Now())) if msg.headerPacket != -1 { t.Error("Should not find headers end") } @@ -183,7 +184,7 @@ func TestTCPMessageMethodType(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if msg.methodType != tc.expectedMethodType { t.Errorf("Expected %d, got %d", tc.expectedMethodType, msg.methodType) @@ -208,7 +209,7 @@ func TestTCPMessageBodyType(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if msg.bodyType != tc.expectedBodyType { t.Errorf("Expected %d, got %d", tc.expectedBodyType, msg.bodyType) @@ -229,12 +230,12 @@ func TestTCPMessageBodySize(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payloads[0]))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payloads[0]), time.Now())) if len(tc.payloads) > 1 { for _, p := range tc.payloads[1:] { seq := uint32(1 + msg.Size()) - msg.AddPacket(buildPacket(tc.direction, 1, seq, []byte(p))) + msg.AddPacket(buildPacket(tc.direction, 1, seq, []byte(p), time.Now())) } } @@ -243,3 +244,15 @@ func TestTCPMessageBodySize(t *testing.T) { } } } + +func TestTcpMessageStart(t *testing.T) { + start := time.Now().Add(-1 * time.Second) + + msg := buildMessage(buildPacket(true, 1, 2, []byte("b"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"), start)) + + if msg.Start != start { + t.Error("Message timestamp should be equal to the lowest related packet timestamp", start, msg.Start) + } +} + diff --git a/raw_socket_listener/tcp_packet.go b/raw_socket_listener/tcp_packet.go index 0e3e832..1dafcbe 100644 --- a/raw_socket_listener/tcp_packet.go +++ b/raw_socket_listener/tcp_packet.go @@ -5,6 +5,7 @@ import ( "log" "strconv" "strings" + "time" ) var _ = log.Println @@ -38,14 +39,16 @@ type TCPPacket struct { Raw []byte Data []byte Addr []byte + timestamp time.Time ID tcpID } // ParseTCPPacket takes address and tcp payload and returns parsed TCPPacket -func ParseTCPPacket(addr []byte, data []byte) (p *TCPPacket) { +func ParseTCPPacket(addr []byte, data []byte, timestamp time.Time) (p *TCPPacket) { p = &TCPPacket{Raw: data} p.ParseBasic() p.Addr = addr + p.timestamp = timestamp p.GenID() return @@ -79,27 +82,33 @@ func (t *TCPPacket) ParseBasic() { t.Data = t.Raw[t.DataOffset*4:] } -func (t *TCPPacket) Dump() []byte { - buf := make([]byte, len(t.Data)+16+16) - copy(buf[:16], t.Addr) +func (t *TCPPacket) Dump() *Packet { - tcpBuf := buf[16:] + packetSrcIP := make([]byte, 16) + packetData := make([]byte, len(t.Data) + 16) - binary.BigEndian.PutUint16(tcpBuf[2:4], t.DestPort) - binary.BigEndian.PutUint16(tcpBuf[0:2], t.SrcPort) + copy(packetSrcIP, t.Addr) - binary.BigEndian.PutUint32(tcpBuf[4:8], t.Seq) - binary.BigEndian.PutUint32(tcpBuf[8:12], t.Ack) + binary.BigEndian.PutUint16(packetData[0:2], t.SrcPort) + binary.BigEndian.PutUint16(packetData[2:4], t.DestPort) - tcpBuf[12] = 64 + binary.BigEndian.PutUint32(packetData[4:8], t.Seq) + binary.BigEndian.PutUint32(packetData[8:12], t.Ack) + + packetData[12] = 64 if t.IsFIN { - tcpBuf[13] = tcpBuf[13] | 0x01 + packetData[13] = packetData[13] | 0x01 } - copy(tcpBuf[16:], t.Data) + copy(packetData[16:], t.Data) + + return &Packet{ + srcIP: packetSrcIP, + data:packetData, + timestamp:t.timestamp, + } - return buf } // String output for a TCP Packet From c6244435caa8f0658905c7fac73564a5e4b53136 Mon Sep 17 00:00:00 2001 From: Yohan Legat Date: Thu, 1 Dec 2016 10:52:21 +0100 Subject: [PATCH 18/20] [Boyscout] unexport TCPPacket.Dump() and Packet struct --- raw_socket_listener/listener.go | 10 +++++----- raw_socket_listener/listener_test.go | 22 +++++++++++----------- raw_socket_listener/tcp_packet.go | 4 ++-- 3 files changed, 18 insertions(+), 18 deletions(-) diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index 1b1c8c0..8428ba7 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -33,7 +33,7 @@ import ( var _ = fmt.Println -type Packet struct { +type packet struct { srcIP []byte data []byte timestamp time.Time @@ -59,7 +59,7 @@ type Listener struct { respWithoutReq map[uint32]tcpID // Messages ready to be send to client - packetsChan chan *Packet + packetsChan chan *packet // Messages ready to be send to client messagesChan chan *TCPMessage @@ -94,7 +94,7 @@ const ( func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration) (l *Listener) { l = &Listener{} - l.packetsChan = make(chan *Packet, 10000) + l.packetsChan = make(chan *packet, 10000) l.messagesChan = make(chan *TCPMessage, 10000) l.quit = make(chan bool) l.readyCh = make(chan bool, 1) @@ -636,14 +636,14 @@ func (t *Listener) readRAWSocket() { } } -func (t *Listener) buildPacket(packetSrcIP []byte, packetData []byte, timestamp time.Time) *Packet { +func (t *Listener) buildPacket(packetSrcIP []byte, packetData []byte, timestamp time.Time) *packet { copyPacketSrcIP := make([]byte, 16) copyPacketData := make([]byte, len(packetData)) copy(copyPacketSrcIP, packetSrcIP) copy(copyPacketData, packetSrcIP) - return &Packet { + return &packet{ srcIP: packetSrcIP, data: packetData, timestamp:timestamp, diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 8ade26e..70271a7 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -20,8 +20,8 @@ func TestRawListenerInput(t *testing.T) { respAck := reqPacket.Seq + uint32(len(reqPacket.Data)) respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) - listener.packetsChan <- reqPacket.Dump() - listener.packetsChan <- respPacket.Dump() + listener.packetsChan <- reqPacket.dump() + listener.packetsChan <- respPacket.dump() select { case req = <-listener.messagesChan: @@ -59,9 +59,9 @@ func TestRawListenerInputResponseByClose(t *testing.T) { finPacket := buildPacket(false, respAck, reqPacket.Seq+2, []byte(""), time.Now()) finPacket.IsFIN = true - listener.packetsChan <- reqPacket.Dump() - listener.packetsChan <- respPacket.Dump() - listener.packetsChan <- finPacket.Dump() + listener.packetsChan <- reqPacket.dump() + listener.packetsChan <- respPacket.dump() + listener.packetsChan <- finPacket.dump() select { case req = <-listener.messagesChan: @@ -94,7 +94,7 @@ func TestRawListenerInputWithoutResponse(t *testing.T) { reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) - listener.packetsChan <- reqPacket.Dump() + listener.packetsChan <- reqPacket.dump() select { case req = <-listener.messagesChan: @@ -118,8 +118,8 @@ func TestRawListenerResponse(t *testing.T) { respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) // If response packet comes before request - listener.packetsChan <- respPacket.Dump() - listener.packetsChan <- reqPacket.Dump() + listener.packetsChan <- respPacket.dump() + listener.packetsChan <- reqPacket.dump() select { case req = <-listener.messagesChan: @@ -209,7 +209,7 @@ func TestAlt100ContinueHeaderOrder(t *testing.T) { func testRawListener100Continue(t *testing.T, listener *Listener, result []byte, packets ...*TCPPacket) { var req, resp *TCPMessage for _, p := range packets { - listener.packetsChan <- p.Dump() + listener.packetsChan <- p.dump() } select { @@ -249,7 +249,7 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket var r, req, resp *TCPMessage for _, p := range packets { - listener.packetsChan <- p.Dump() + listener.packetsChan <- p.dump() } select { @@ -452,7 +452,7 @@ func TestRawListenerBench(t *testing.T) { } } - l.packetsChan <- p.Dump() + l.packetsChan <- p.dump() time.Sleep(time.Millisecond) } diff --git a/raw_socket_listener/tcp_packet.go b/raw_socket_listener/tcp_packet.go index 1dafcbe..f2ba82c 100644 --- a/raw_socket_listener/tcp_packet.go +++ b/raw_socket_listener/tcp_packet.go @@ -82,7 +82,7 @@ func (t *TCPPacket) ParseBasic() { t.Data = t.Raw[t.DataOffset*4:] } -func (t *TCPPacket) Dump() *Packet { +func (t *TCPPacket) dump() *packet { packetSrcIP := make([]byte, 16) packetData := make([]byte, len(t.Data) + 16) @@ -103,7 +103,7 @@ func (t *TCPPacket) Dump() *Packet { copy(packetData[16:], t.Data) - return &Packet{ + return &packet{ srcIP: packetSrcIP, data:packetData, timestamp:t.timestamp, From b462c06a7c2f28ec08e88f504992563c00e80c89 Mon Sep 17 00:00:00 2001 From: Yohan Legat Date: Thu, 8 Dec 2016 17:35:45 +0100 Subject: [PATCH 19/20] remove dead code --- raw_socket_listener/listener.go | 6 ------ 1 file changed, 6 deletions(-) diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index 8428ba7..9a8b83f 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -528,12 +528,6 @@ func (t *Listener) readPcap() { } } - packetSrcIP := make([]byte, 16) - packetData := make([]byte, len(data)) - - copy(packetSrcIP, srcIP) - copy(packetData, data) - t.packetsChan <- t.buildPacket(srcIP, data, packet.Metadata().Timestamp) } } From 0a74db6aae5c0c4dba0af0622637a6a87c9a51f3 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Sun, 18 Dec 2016 20:09:24 +0300 Subject: [PATCH 20/20] Disallow zero timeout (also used fix tests) --- http_client.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/http_client.go b/http_client.go index 8007055..aeb322f 100644 --- a/http_client.go +++ b/http_client.go @@ -59,6 +59,10 @@ func NewHTTPClient(baseURL string, config *HTTPClientConfig) *HTTPClient { u, _ := url.Parse(baseURL) + if config.Timeout == 0 { + config.Timeout = time.Second + } + config.ConnectionTimeout = config.Timeout if config.ResponseBufferSize == 0 {