From 4255d4155fb7d806e139d94e62e74db1086bf77f Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Fri, 15 Sep 2017 19:25:46 +0500 Subject: [PATCH 01/37] Improve handling for 100-continue requests for clients ignoring response (#513) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Some clients, even if have `Expect: 100-continue` in request, start sending data without waiting approval from the server. More over, there is cases when response “100 Continue” is received, client changes TCP seq, and following requests start having valid sequences, based on response one. This change should address this bugs. As small addition extended list of http methods, to fully support webdav. --- proto/proto.go | 2 +- raw_socket_listener/listener.go | 27 ++-- raw_socket_listener/listener_test.go | 206 +++++++++++++-------------- raw_socket_listener/tcp_message.go | 9 ++ 4 files changed, 127 insertions(+), 117 deletions(-) diff --git a/proto/proto.go b/proto/proto.go index d149d0f..107de3f 100644 --- a/proto/proto.go +++ b/proto/proto.go @@ -467,7 +467,7 @@ func Status(payload []byte) []byte { } var httpMethods []string = []string{ - "GET ", "OPTI", "HEAD", "POST", "PUT ", "DELE", "TRAC", "CONN", "PATC" /* custom methods */, "BAN", "PURG", + "GET ", "OPTI", "HEAD", "POST", "PUT ", "DELE", "TRAC", "CONN", "PATC" /* custom methods */, "BAN ", "PURG", "PROP", "MKCO", "COPY", "MOVE", "LOCK", "UNLO", } func IsHTTPPayload(payload []byte) bool { diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index d8a1478..518cc06 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -693,12 +693,20 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { } }() + // log.Println("PACKET:", packet, t.seqWithData) + + var responseRequest *TCPMessage var message *TCPMessage isIncoming := packet.DestPort == t.port + if !isIncoming { + responseRequest, _ = t.respAliases[packet.Ack] + } + // Seek for 100-expect chunks - if parentAck, ok := t.seqWithData[packet.Seq]; ok { + // `packet.Ack != parentAck` is protection for clients who send data without ignoring server 100-continue response, e.g have data chunks have same Ack + if parentAck, ok := t.seqWithData[packet.Seq]; ok && packet.Ack != parentAck { // Skip zero-length chunks https://github.com/buger/goreplay/issues/496 if len(packet.Data) == 0 { return @@ -739,12 +747,6 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { packet.UpdateAck(alias) } - var responseRequest *TCPMessage - - if !isIncoming { - responseRequest, _ = t.respAliases[packet.Ack] - } - message, ok := t.messages[packet.ID] if !ok { @@ -766,8 +768,9 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { // Handling Expect: 100-continue requests if message.expectType == httpExpect100Continue && len(message.packets) == message.headerPacket+1 { - seq := packet.Seq + uint32(message.Size()) + seq := packet.Seq + uint32(len(packet.Data)) t.seqWithData[seq] = packet.Ack + message.DataSeq = seq message.complete = false @@ -793,7 +796,13 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { packet.Data = proto.DeleteHeader(packet.Data, bExpectHeader) } - // log.Println("Received message:", string(message.Bytes()), message.ID(), t.messages) + // If client do sends Expect: 100-continue but do not respect server response + if message.expectType == httpExpect100Continue && (message.headerPacket != -1 && len(message.packets) > message.headerPacket+1) { + delete(t.seqWithData, message.DataSeq) + seq := packet.Seq + uint32(len(packet.Data)) + t.seqWithData[seq] = packet.Ack + message.DataSeq = seq + } if isIncoming { // If message have multiple packets, delete previous alias diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index dbacdcd..988cd65 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -46,79 +46,100 @@ func TestRawListenerInput(t *testing.T) { } } +func firstPacket(payload []byte) *TCPPacket { + return buildPacket( + true, + 1, + 1, + payload, + time.Now(), + ) +} + +func nextPacket(prev *TCPPacket, payload []byte) *TCPPacket { + return buildPacket( + prev.SrcPort == 1, + prev.Ack, + prev.Seq+uint32(len(prev.Data)), + payload, + time.Now(), + ) +} + +func responsePacket(prev *TCPPacket, payload []byte) *TCPPacket { + return buildPacket( + !(prev.SrcPort == 1), + prev.Seq+uint32(len(prev.Data)), + prev.Ack, + payload, + time.Now(), + ) +} + func TestSingleAck100Continue(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: 4\r\n\r\n"), time.Now()) - - respPacket1 := buildPacket(false, - uint32(len(reqPacket1.Data)) + reqPacket1.Seq, - 1, - []byte(""), time.Now()) - - respPacket2 := buildPacket( false, - uint32(len(reqPacket1.Data)) + reqPacket1.Seq, - 1, - []byte("HTTP/1.1 100 Continue\r\n"), time.Now()) - - reqPacket3 := buildPacket(true, - uint32(len(reqPacket1.Data)) + respPacket1.Seq, - reqPacket1.Seq+uint32(len(reqPacket1.Data)), - []byte("DATA"), time.Now()) - - respPacket3 := buildPacket(false, - uint32(len(reqPacket3.Data)) + reqPacket3.Seq, - respPacket1.Seq+uint32(len(respPacket1.Data)), []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) + reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) + respPacket1 := responsePacket(reqPacket1, []byte("")) + respPacket2 := responsePacket(reqPacket1, []byte("HTTP/1.1 100 Continue\r\n")) + reqPacket2 := responsePacket(respPacket2, []byte("DATA")) + respPacket3 := responsePacket(reqPacket2, []byte("HTTP/1.1 200 OK\r\n\r\n")) result := []byte("POST / HTTP/1.1\r\nContent-Length: 4\r\n\r\nDATA") testRawListener100Continue(t, listener, result, reqPacket1, respPacket1, respPacket2, - reqPacket3, - respPacket3 ) + reqPacket2, + respPacket3) } +func Test100ContinueWithoutWaiting(t *testing.T) { + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + defer listener.Close() + + req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) + req2 := nextPacket(req1, []byte("DATA")) + resp1 := responsePacket(req1, []byte("HTTP/1.1 100 Continue\r\n")) + resp2 := responsePacket(req2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + + result := []byte("POST / HTTP/1.1\r\nContent-Length: 4\r\n\r\nDATA") + + testRawListener100Continue(t, listener, result, + req1, req2, resp1, resp2) +} + +// Client first sends data without waiting 100-continue, but once response received, generate packets based on Ack payload +func Test100ContinueMixed(t *testing.T) { + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + defer listener.Close() + + req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 12\r\n\r\n")) + req2 := nextPacket(req1, []byte("DAT1")) + resp1 := responsePacket(req1, []byte("HTTP/1.1 100 Continue\r\n\r\n")) + req3 := responsePacket(resp1, []byte("DAT2")) + req3.Seq = req2.Seq + uint32(len(req2.Data)) + req4 := nextPacket(req3, []byte("DAT3")) + resp2 := responsePacket(req4, []byte("HTTP/1.1 200 OK\r\n\r\n")) + + result := []byte("POST / HTTP/1.1\r\nContent-Length: 12\r\n\r\nDAT1DAT2DAT3") + + testRawListener100Continue(t, listener, result, + req1, req2, req3, req4, resp1, resp2) +} func TestDoubleAck100Continue(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: 4\r\n\r\n"), time.Now()) + reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) - respPacket1 := buildPacket(false, - uint32(len(reqPacket1.Data)) + reqPacket1.Seq, - 1, - []byte(""), time.Now()) - - respPacket2 := buildPacket( false, - uint32(len(reqPacket1.Data)) + reqPacket1.Seq, - 1, - []byte("HTTP/1.1 100 Continue\r\n"), time.Now()) - - reqPacket2 := buildPacket(true, - uint32(len(reqPacket1.Data)) + respPacket1.Seq, - reqPacket1.Seq+uint32(len(reqPacket1.Data)), - []byte(""), time.Now()) - - reqPacket3 := buildPacket(true, - uint32(len(reqPacket1.Data)) + respPacket1.Seq, - reqPacket1.Seq+uint32(len(reqPacket1.Data)), - []byte("DATA"), time.Now()) - - respPacket3 := buildPacket(false, - uint32(len(reqPacket3.Data)) + reqPacket3.Seq, - respPacket1.Seq+uint32(len(respPacket1.Data)), - []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) + respPacket1 := responsePacket(reqPacket1, []byte("")) + respPacket2 := responsePacket(reqPacket1, []byte("HTTP/1.1 100 Continue\r\n")) + reqPacket2 := responsePacket(respPacket2, []byte("")) + reqPacket3 := responsePacket(respPacket2, []byte("DATA")) + respPacket3 := responsePacket(reqPacket3, []byte("HTTP/1.1 200 OK\r\n\r\n")) result := []byte("POST / HTTP/1.1\r\nContent-Length: 4\r\n\r\nDATA") @@ -126,10 +147,9 @@ func TestDoubleAck100Continue(t *testing.T) { reqPacket1, respPacket1, respPacket2, reqPacket2, reqPacket3, - respPacket3 ) + respPacket3) } - func TestRawListenerInputResponseByClose(t *testing.T) { var req, resp *TCPMessage @@ -198,8 +218,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"), time.Now()) - respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) + reqPacket := firstPacket([]byte("GET / HTTP/1.1\r\n\r\n")) + respPacket := responsePacket(reqPacket, []byte("HTTP/1.1 200 OK\r\n\r\n")) // If response packet comes before request listener.packetsChan <- respPacket.dump() @@ -232,23 +252,25 @@ func TestRawListenerResponse(t *testing.T) { } } +func get100ContinuePackets() (req []*TCPPacket, resp []*TCPPacket) { + req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 2\r\n\r\n")) + resp1 := responsePacket(req1, []byte("HTTP/1.1 100 Continue\r\n")) + req2 := responsePacket(resp1, []byte("a")) + req3 := nextPacket(req2, []byte("b")) + resp2 := responsePacket(req3, []byte("HTTP/1.1 200 OK\r\n\r\n")) + + return []*TCPPacket{req1, req2, req3}, []*TCPPacket{resp1, resp2} +} + 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"), time.Now()) - // Packet with data have different Seq - 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"), 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"), time.Now()) + req, resp := get100ContinuePackets() result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") - testRawListener100Continue(t, listener, result, reqPacket1, reqPacket2, reqPacket3, respPacket1, respPacket2) + testRawListener100Continue(t, listener, result, req[0], req[1], req[2], resp[0], resp[1]) } // Response comes before Request @@ -256,38 +278,11 @@ 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"), time.Now()) - // Packet with data have different Seq - 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"), 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"), time.Now()) + req, resp := get100ContinuePackets() result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") - testRawListener100Continue(t, listener, result, respPacket1, respPacket2, reqPacket1, reqPacket2, reqPacket3) -} - -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"), time.Now()) - // Packet with data have different Seq - 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"), 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"), time.Now()) - - result := []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\nab") - - testRawListener100Continue(t, listener, result, reqPacket1, reqPacket2, reqPacket3, respPacket1, respPacket2) + testRawListener100Continue(t, listener, result, resp[0], resp[1], req[0], req[1], req[2]) } func testRawListener100Continue(t *testing.T, listener *Listener, result []byte, packets ...*TCPPacket) { @@ -436,22 +431,19 @@ 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"), time.Now()) - // Packet with data have different Seq - 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()) + reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n")) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n"), time.Now()) + respPacket1 := responsePacket(reqPacket1, []byte("HTTP/1.1 100 Continue\r\n")) + reqPacket2 := responsePacket(respPacket1, []byte("1\r\na\r\n")) + reqPacket3 := nextPacket(reqPacket2, []byte("1\r\nb\r\n")) + reqPacket4 := nextPacket(reqPacket3, []byte("0\r\n\r\n")) - // 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"), time.Now()) + respPacket2 := responsePacket(reqPacket4, []byte("HTTP/1.1 200 OK\r\n\r\n")) // Should re-construct message from all possible combinations for i := 0; i < 6*5*4*3*2*1; i++ { packets := permutation(i, []*TCPPacket{reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket1, respPacket2}) - t.Log("permutation:", i) testChunkedSequence(t, listener, packets...) } } diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 02f8384..86bc8b0 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -8,6 +8,7 @@ import ( "log" "net" "strconv" + "strings" "time" "github.com/buger/goreplay/proto" @@ -479,3 +480,11 @@ func (t *TCPMessage) ID() tcpID { func (t *TCPMessage) IP() net.IP { return net.IP(t.packets[0].Addr) } + +func (t *TCPMessage) String() string { + return strings.Join([]string{ + "Len packets: " + strconv.Itoa(len(t.packets)), + "Data size:" + strconv.Itoa(len(t.Bytes())), + "Data:" + string(t.Bytes()), + }, "\n") +} From c5f12b567bad4614feb37f0f5d7c2f409292bf18 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Sun, 24 Sep 2017 22:32:03 +0500 Subject: [PATCH 02/37] Add `log` helper and fix memory leak --- middleware/middleware.js | 11 ++++++++++- middleware/package.json | 2 +- 2 files changed, 11 insertions(+), 2 deletions(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index 306c17b..d50ee3e 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -41,12 +41,17 @@ function init() { let resp = msg; - ["message", chanPrefix, chanPrefix + "#" + msg.ID].forEach(function(chanID){ + ["message", chanPrefix, chanPrefix + "#" + msg.ID].forEach(function(chanID, idx){ if (proxy.ch[chanID]) { proxy.ch[chanID].forEach(function(ch){ let r = ch.cb(msg); if (r) resp = r; // If one of callback decided not to send response back, do not override it in global callbacks }) + + // Cleanup Individual message channels to avoid memory leaks + if (idx == 2) { + delete proxy.ch[chanID] + } } }) @@ -389,6 +394,10 @@ function fail(message) { console.error("\x1b[31m[MIDDLEWARE] %s\x1b[0m", message) } +function log(message) { + console.error(message) +} + function TEST_init() { const child_process = require('child_process'); diff --git a/middleware/package.json b/middleware/package.json index 698165f..007ed27 100644 --- a/middleware/package.json +++ b/middleware/package.json @@ -1,6 +1,6 @@ { "name": "goreplay_middleware", - "version": "0.1.13", + "version": "0.1.14", "description": "Package for writing middleware for GoReplay https://goreplay.org", "main": "middleware.js", "scripts": { From debb265d6d9a8393aa24d850137351329d6a5bf1 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Wed, 4 Oct 2017 20:18:05 +0300 Subject: [PATCH 03/37] Fix memory leak --- middleware/middleware.js | 54 +++++++++++++++++++++++++++++++++++++--- middleware/package.json | 5 ++-- 2 files changed, 53 insertions(+), 6 deletions(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index d50ee3e..2548c7b 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -62,14 +62,25 @@ function init() { } // Clean up old messaged ID specific channels if they are older then 60s - setInterval(function(){ + let gc = function(gcTime){ let now = new Date(); for (k in proxy.ch) { if (k.indexOf("#") == -1) continue; - proxy.ch[k] = proxy.ch[k].filter(function(ch){ return (now - ch.created) < (60 * 1000) }) + proxy.ch[k] = proxy.ch[k].filter(function(ch){ + return (now - ch.created) < gcTime + }) + + if (proxy.ch[k].length == 0) { + delete proxy.ch[k] + } } - }, 1000) + } + proxy.gc = gc + + setInterval(function(){ + gc(10 * 1000) + }, 1000); const readline = require('readline'); const rl = readline.createInterface({ @@ -375,7 +386,8 @@ module.exports = { setHttpBodyParam: setHttpBodyParam, httpCookie: httpCookie, setHttpCookie: setHttpCookie, - test: testRunner + test: testRunner, + benchmark: testBenchmark } @@ -389,6 +401,40 @@ function testRunner(){ }) } +function testBenchmark(){ + const child_process = require('child_process'); + + let gor = init(); + gor.on("message", function(){ + }); + + gor.on("request", function(){ + }); + + for (var i = 0; i<256; i++) { + let req = parseMessage(Buffer.from("1 2 3\nGET / HTTP/1.1\r\n\r\n").toString('hex')); + req.ID = +Date.now() + gor.emit(req); + + gor.on("request", req.ID+"", function(){ + gor.on("response", req.ID+"", function(){ + }) + }) + + if ( i % 3 == 0 ) { + let resp = parseMessage(Buffer.from("2 2 3\nHTTP/1.1 200 OK\r\n\r\n").toString('hex')); + resp.ID = req.ID + gor.emit(resp); + } + } + + child_process.execSync("sleep 0.01"); + + gor.gc(1) + + fail(JSON.stringify(gor.ch)) +} + // Just print in red color function fail(message) { console.error("\x1b[31m[MIDDLEWARE] %s\x1b[0m", message) diff --git a/middleware/package.json b/middleware/package.json index 007ed27..9748460 100644 --- a/middleware/package.json +++ b/middleware/package.json @@ -1,10 +1,11 @@ { "name": "goreplay_middleware", - "version": "0.1.14", + "version": "0.1.15", "description": "Package for writing middleware for GoReplay https://goreplay.org", "main": "middleware.js", "scripts": { - "test": "node -e \"var gor = require('./middleware.js'); gor.test(); process.exit()\"" + "test": "node -e \"var gor = require('./middleware.js'); gor.test(); process.exit()\"", + "benchmark": "node -e \"var gor = require('./middleware.js'); gor.benchmark(); process.exit()\"" }, "keywords": [ "middleware", From cc005ddcbbb0c033d706721ca2de03621eedbb89 Mon Sep 17 00:00:00 2001 From: noibar Date: Mon, 23 Oct 2017 18:26:14 +0300 Subject: [PATCH 04/37] Fix comment on tcp message methods (#521) * Fix comment on tcp message methods * Update tcp message constructor documentation --- raw_socket_listener/tcp_message.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 86bc8b0..839a47a 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -49,7 +49,8 @@ type TCPMessage struct { complete bool } -// NewTCPMessage pointer created from a Acknowledgment number and a channel of messages readuy to be deleted +// NewTCPMessage pointer created from a sequence and acknowledgment numbers, whether the message is incoming and a timestamp +// that indicates when the packet was captrued. func NewTCPMessage(Seq, Ack uint32, IsIncoming bool, timestamp time.Time) (msg *TCPMessage) { msg = &TCPMessage{Seq: Seq, Ack: Ack, IsIncoming: IsIncoming, Start: timestamp} @@ -74,7 +75,7 @@ func (t *TCPMessage) Bytes() (output []byte) { return output } -// Size returns total body size +// BodySize returns total body size func (t *TCPMessage) BodySize() (size int) { if len(t.packets) == 0 || t.headerPacket == -1 { return 0 @@ -225,7 +226,7 @@ func (t *TCPMessage) updateHeadersPacket() { return } -// isMultipart returns true if message contains from multiple tcp packets +// checkIfComplete returns true if all of the packets that compse the message arrived. func (t *TCPMessage) checkIfComplete() { if t.seqMissing || t.headerPacket == -1 { return From 3acec0fd4239287c6c5983f33bdc086f34a3b1ee Mon Sep 17 00:00:00 2001 From: Mel Shafer Date: Mon, 30 Oct 2017 22:15:52 -0400 Subject: [PATCH 05/37] fix typo in the README --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 0b09462..8cd1283 100644 --- a/README.md +++ b/README.md @@ -43,7 +43,7 @@ We have created a [GoReplay PRO](https://goreplay.org/pro.html) extension which ## Problems? If you have a problem, please review the [FAQ](https://github.com/buger/goreplay/wiki/FAQ) and [Troubleshooting](https://github.com/buger/goreplay/wiki/Troubleshooting) wiki pages. Searching the [issues](https://github.com/buger/goreplay/issues) for your problem is also a good idea. -All bug-reports and suggestions should go though Github Issues or our [Google Group](https://groups.google.com/forum/#!forum/gor-users) (you can just send email to gor-users@googlegroups.com). +All bug-reports and suggestions should go through Github Issues or our [Google Group](https://groups.google.com/forum/#!forum/gor-users) (you can just send email to gor-users@googlegroups.com). If you have a private question feel free to send email to support@gortool.com. From bf8198d6391b4713b39870be6fd8e17910ee3af9 Mon Sep 17 00:00:00 2001 From: Daniel Albuquerque Date: Tue, 14 Nov 2017 22:13:39 +0000 Subject: [PATCH 06/37] Hardcode alpine version and bump version of goreplay from 16.0.2 to 16.1 --- Dockerfile | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/Dockerfile b/Dockerfile index 914094e..3f03461 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,5 +1,5 @@ -FROM alpine:latest +FROM alpine:3.6 RUN apk update && apk add ca-certificates && update-ca-certificates && apk add openssl -RUN wget https://github.com/buger/goreplay/releases/download/v0.16.0.2/gor_0.16.0_x64.tar.gz -O gor.tar.gz +RUN wget https://github.com/buger/goreplay/releases/download/v0.16.1/gor_0.16.1_x64.tar.gz -O gor.tar.gz RUN tar xzf gor.tar.gz -ENTRYPOINT ./gor +ENTRYPOINT ["./goreplay"] \ No newline at end of file From 45d7527128df3170911c6a90688c0ec0da3021a3 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Wed, 10 Jan 2018 21:23:55 +0200 Subject: [PATCH 07/37] Increase middleware buffer to avoid overloads when prettifier returns buffer larger then original --- middleware.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/middleware.go b/middleware.go index 83ad8d1..0b720e0 100644 --- a/middleware.go +++ b/middleware.go @@ -62,7 +62,7 @@ func (m *Middleware) ReadFrom(plugin io.Reader) { func (m *Middleware) copy(to io.Writer, from io.Reader) { buf := make([]byte, 5*1024*1024) - dst := make([]byte, len(buf)*2) + dst := make([]byte, len(buf)*3) for { nr, _ := from.Read(buf) From 22a67df6c0a0ef0d3639d6d3bbcf85a34d22bd80 Mon Sep 17 00:00:00 2001 From: Ohad Basan Date: Wed, 10 Jan 2018 18:48:44 +0200 Subject: [PATCH 08/37] fix wrong param in docs --- docs/Request-rewriting.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/Request-rewriting.md b/docs/Request-rewriting.md index 825610b..fb378b7 100644 --- a/docs/Request-rewriting.md +++ b/docs/Request-rewriting.md @@ -23,16 +23,16 @@ Set request header, if header already exists it will be overwritten. May be usef ``` gor --input-raw :80 --output-http "http://staging.server" \ - --http-header "User-Agent: Replayed by Gor" \ - --http-header "Enable-Feature-X: true" + --http-set-header "User-Agent: Replayed by Gor" \ + --http-set-header "Enable-Feature-X: true" ``` #### Host header -Host header gets special treatment. By default Host get set to the value specified in --output-http. If you manually set --http-header "Host: anonther.com", Gor will not override Host value. +Host header gets special treatment. By default Host get set to the value specified in --output-http. If you manually set --http-set-header "Host: anonther.com", Gor will not override Host value. If you app accepts traffic from multiple domains, and you want to keep original headers, there is specific `--http-original-host` with tells Gor do not touch Host header at all. *** -You may also read about [[Request filtering]], [[Rate limiting]] and [[Middleware]] \ No newline at end of file +You may also read about [[Request filtering]], [[Rate limiting]] and [[Middleware]] From f8a1b9dde6dc6d51630ef33d774db7a2d2e44f8b Mon Sep 17 00:00:00 2001 From: Ohad Basan Date: Thu, 11 Jan 2018 16:15:52 +0200 Subject: [PATCH 09/37] fix error in bash example --- docs/Middleware.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/Middleware.md b/docs/Middleware.md index 7c98cbe..5edd111 100644 --- a/docs/Middleware.md +++ b/docs/Middleware.md @@ -26,7 +26,7 @@ Simple bash echo middleware (returns same request) will look like this: ```bash while read line; do echo $line -end +done ``` Middleware can be enabled using `--middleware` option, by specifying path to executable file: @@ -71,4 +71,4 @@ Imagine that you have auth system that randomly generate access tokens, which us *** -You may also read about [[Request filtering]], [[Rate limiting]] and [[Request rewriting]]. \ No newline at end of file +You may also read about [[Request filtering]], [[Rate limiting]] and [[Request rewriting]]. From 0097c599df9bd4f686676123e87b37950f47d831 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Fri, 2 Feb 2018 20:45:25 +0200 Subject: [PATCH 10/37] Fix response latency calculation Previously latency calcualted as Response.End - Request.Start Where both End and Start is a last and first packets This calcualtion is wrong, because it is total roundtrip Correct server latency will be Response.End - Request.End In addition added new `--input-raw-timestamp-type` option which allows choose more precise packet timestamp source (if available). --- elasticsearch.go | 6 +++--- elasticsearch_test.go | 12 ++++++------ http_client.go | 12 ++++++------ http_modifier.go | 4 ++-- http_modifier_settings.go | 1 - http_modifier_test.go | 4 ++-- input_raw.go | 8 +++++--- input_raw_test.go | 14 +++++++------- middleware_test.go | 4 ++-- output_file.go | 4 ++-- plugins.go | 2 +- raw_socket_listener/listener.go | 27 ++++++++++++++++++++++++--- raw_socket_listener/tcp_message.go | 12 +++++------- settings.go | 7 +++++-- 14 files changed, 70 insertions(+), 47 deletions(-) diff --git a/elasticsearch.go b/elasticsearch.go index b83dabf..7a8bd3f 100644 --- a/elasticsearch.go +++ b/elasticsearch.go @@ -1,9 +1,9 @@ package main import ( - "net/url" "encoding/json" "log" + "net/url" "strings" //"regexp" "time" @@ -72,10 +72,10 @@ func parseURI(URI string) (err error, index string) { // check URL validity by extracting host and undex values. host := parsedUrl.Host urlPathParts := strings.Split(parsedUrl.Path, "/") - index = urlPathParts[len(urlPathParts) - 1 ] + index = urlPathParts[len(urlPathParts)-1] // force index specification in uri : ie no implicit index - if (host == "" || index == "") { + if host == "" || index == "" { err = new(ESUriErorr) } diff --git a/elasticsearch_test.go b/elasticsearch_test.go index 701c38f..27e624b 100644 --- a/elasticsearch_test.go +++ b/elasticsearch_test.go @@ -6,19 +6,19 @@ import ( const expectedIndex = "gor" -func assertExpectedGorIndex (index string, t *testing.T) { +func assertExpectedGorIndex(index string, t *testing.T) { if expectedIndex != index { t.Fatalf("Expected index %s but got %s", expectedIndex, index) } } -func assertExpectedIndex (expectedIndex string, index string, t *testing.T) { +func assertExpectedIndex(expectedIndex string, index string, t *testing.T) { if expectedIndex != index { t.Fatalf("Expected index %s but got %s", expectedIndex, index) } } -func assertExpectedError (returnedError error, t *testing.T) { +func assertExpectedError(returnedError error, t *testing.T) { expectedError := new(ESUriErorr) if expectedError != returnedError { @@ -26,7 +26,7 @@ func assertExpectedError (returnedError error, t *testing.T) { } } -func assertNoError (returnedError error, t *testing.T) { +func assertNoError(returnedError error, t *testing.T) { if nil != returnedError { t.Errorf("Expected err %s but got %s", nil, returnedError) } @@ -38,7 +38,7 @@ func assertNoError (returnedError error, t *testing.T) { func TestElasticConnectionBuildFailWithoutScheme(t *testing.T) { uri := "localhost:9200/" + expectedIndex - err, _ := parseURI(uri) + err, _ := parseURI(uri) assertExpectedError(err, t) } @@ -48,7 +48,7 @@ func TestElasticConnectionBuildFailWithoutScheme(t *testing.T) { func TestElasticConnectionBuildFailWithoutIndex(t *testing.T) { uri := "http://localhost:9200" - err, index := parseURI(uri) + err, index := parseURI(uri) assertExpectedIndex("", index, t) diff --git a/http_client.go b/http_client.go index 6bc2c11..c264e4b 100644 --- a/http_client.go +++ b/http_client.go @@ -298,9 +298,9 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { if err != nil && readBytes == 0 { maxRead := 100 - if readBytes < maxRead { - maxRead = readBytes - } + if readBytes < maxRead { + maxRead = readBytes + } Debug("[HTTPClient] Response read timeout error", err, c.conn, readBytes, string(c.respBuf[:maxRead])) response = errorPayload(HTTP_TIMEOUT) c.Disconnect() @@ -309,9 +309,9 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { if readBytes < 4 || string(c.respBuf[:4]) != "HTTP" { maxRead := 100 - if readBytes < maxRead { - maxRead = readBytes - } + if readBytes < maxRead { + maxRead = readBytes + } Debug("[HTTPClient] Response read unknown error", err, c.conn, readBytes, string(c.respBuf[:maxRead])) response = errorPayload(HTTP_UNKNOWN_ERROR) c.Disconnect() diff --git a/http_modifier.go b/http_modifier.go index 7e90208..8ed611e 100644 --- a/http_modifier.go +++ b/http_modifier.go @@ -2,9 +2,9 @@ package main import ( "bytes" - "strings" - "hash/fnv" "encoding/base64" + "hash/fnv" + "strings" "github.com/buger/goreplay/proto" ) diff --git a/http_modifier_settings.go b/http_modifier_settings.go index 5ae70e6..7291162 100644 --- a/http_modifier_settings.go +++ b/http_modifier_settings.go @@ -81,7 +81,6 @@ func (h *HTTPHeaderBasicAuthFilters) Set(value string) error { return nil } - // // Handling of --http-allow-header-hash and --http-allow-param-hash options // diff --git a/http_modifier_test.go b/http_modifier_test.go index e39e0af..ac7baba 100644 --- a/http_modifier_test.go +++ b/http_modifier_test.go @@ -97,7 +97,7 @@ func TestHTTPHeaderBasicAuthFilters(t *testing.T) { payload = []byte("POST /post HTTP/1.1\r\nContent-Length: 88\r\nAuthorization: Basic Y3VzdG9tZXI2OnJlc3RAMTIzXlRFU1Q==\r\n\r\na=1&b=2") if len(modifier.Rewrite(payload)) == 0 { t.Error("Request should pass filters") - } + } filters = HTTPHeaderBasicAuthFilters{} // Setting filter that not match our header @@ -115,7 +115,7 @@ func TestHTTPHeaderBasicAuthFilters(t *testing.T) { payload = []byte("POST /post HTTP/1.1\r\nContent-Length: 88\r\nAuthorization: Basic bWlja2V5IG1vdXNlOmhhcHB5MTIz\r\n\r\na=1&b=2") if len(modifier.Rewrite(payload)) == 0 { t.Error("Request should pass filters") - } + } } func TestHTTPModifierURLRewrite(t *testing.T) { diff --git a/input_raw.go b/input_raw.go index eec5002..7dddce1 100644 --- a/input_raw.go +++ b/input_raw.go @@ -20,6 +20,7 @@ type RAWInput struct { trackResponse bool listener *raw.Listener bpfFilter string + timestampType string } // Available engines for intercepting traffic @@ -30,7 +31,7 @@ const ( ) // NewRAWInput constructor for RAWInput. Accepts address with port as argument. -func NewRAWInput(address string, engine int, trackResponse bool, expire time.Duration, realIPHeader string, bpfFilter string) (i *RAWInput) { +func NewRAWInput(address string, engine int, trackResponse bool, expire time.Duration, realIPHeader string, bpfFilter string, timestampType string) (i *RAWInput) { i = new(RAWInput) i.data = make(chan *raw.TCPMessage) i.address = address @@ -40,6 +41,7 @@ func NewRAWInput(address string, engine int, trackResponse bool, expire time.Dur i.realIPHeader = []byte(realIPHeader) i.quit = make(chan bool) i.trackResponse = trackResponse + i.timestampType = timestampType i.listen(address) i.listener.IsReady() @@ -59,7 +61,7 @@ func (i *RAWInput) Read(data []byte) (int, error) { buf = proto.SetHeader(buf, i.realIPHeader, []byte(msg.IP().String())) } } else { - header = payloadHeader(ResponsePayload, msg.UUID(), msg.AssocMessage.Start.UnixNano(), msg.End.UnixNano()-msg.AssocMessage.Start.UnixNano()) + header = payloadHeader(ResponsePayload, msg.UUID(), msg.Start.UnixNano(), msg.End.UnixNano()-msg.AssocMessage.End.UnixNano()) } copy(data[0:len(header)], header) @@ -77,7 +79,7 @@ func (i *RAWInput) listen(address string) { log.Fatal("input-raw: error while parsing address", err) } - i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.bpfFilter) + i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.bpfFilter, i.timestampType) ch := i.listener.Receiver() diff --git a/input_raw_test.go b/input_raw_test.go index 099b393..cd0b319 100644 --- a/input_raw_test.go +++ b/input_raw_test.go @@ -44,7 +44,7 @@ func TestRAWInputIPv4(t *testing.T) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "X-Real-IP", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "X-Real-IP", "", "") defer input.Close() output := NewTestOutput(func(data []byte) { @@ -106,7 +106,7 @@ func TestRAWInputNoKeepAlive(t *testing.T) { originAddr := listener.Addr().String() - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() output := NewTestOutput(func(data []byte) { @@ -152,7 +152,7 @@ func TestRAWInputIPv6(t *testing.T) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() output := NewTestOutput(func(data []byte) { @@ -203,7 +203,7 @@ func TestInputRAW100Expect(t *testing.T) { originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "") defer input.Close() // We will use it to get content of raw HTTP request @@ -266,7 +266,7 @@ func TestInputRAWChunkedEncoding(t *testing.T) { })) originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "") defer input.Close() replay := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -330,7 +330,7 @@ func TestInputRAWLargePayload(t *testing.T) { })) originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() replay := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { @@ -375,7 +375,7 @@ func BenchmarkRAWInput(b *testing.B) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() output := NewTestOutput(func(data []byte) { diff --git a/middleware_test.go b/middleware_test.go index c828b85..b8d3c9a 100644 --- a/middleware_test.go +++ b/middleware_test.go @@ -118,7 +118,7 @@ func TestEchoMiddleware(t *testing.T) { // Catch traffic from one service fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() // And redirect to another @@ -180,7 +180,7 @@ func TestTokenMiddleware(t *testing.T) { fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) // Catch traffic from one service - input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "") + input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "") defer input.Close() // And redirect to another diff --git a/output_file.go b/output_file.go index 07471f6..0145344 100644 --- a/output_file.go +++ b/output_file.go @@ -186,9 +186,9 @@ func (o *FileOutput) Write(data []byte) (n int, err error) { o.currentID = meta[1] o.payloadType = meta[0] } - + o.updateName() - + if o.file == nil || o.currentName != o.file.Name() { o.mu.Lock() o.Close() diff --git a/plugins.go b/plugins.go index 504abb7..1cb0637 100644 --- a/plugins.go +++ b/plugins.go @@ -106,7 +106,7 @@ func InitPlugins() { } for _, options := range Settings.inputRAW { - registerPlugin(NewRAWInput, options, engine, Settings.inputRAWTrackResponse, Settings.inputRAWExpire, Settings.inputRAWRealIPHeader, Settings.inputRAWBpfFilter) + registerPlugin(NewRAWInput, options, engine, Settings.inputRAWTrackResponse, Settings.inputRAWExpire, Settings.inputRAWRealIPHeader, Settings.inputRAWBpfFilter, Settings.inputRAWTimestampType) } for _, options := range Settings.inputTCP { diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index 518cc06..47a1aba 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -72,7 +72,8 @@ type Listener struct { trackResponse bool messageExpire time.Duration - bpfFilter string + bpfFilter string + timestampType string conn net.PacketConn pcapHandles []*pcap.Handle @@ -95,7 +96,7 @@ const ( ) // NewListener creates and initializes new Listener object -func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration, bpfFilter string) (l *Listener) { +func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration, bpfFilter string, timestampType string) (l *Listener) { l = &Listener{} l.packetsChan = make(chan *packet, 10000) @@ -110,6 +111,7 @@ func NewListener(addr string, port string, engine int, trackResponse bool, expir l.respWithoutReq = make(map[uint32]tcpID) l.trackResponse = trackResponse l.bpfFilter = bpfFilter + l.timestampType = timestampType l.addr = addr _port, _ := strconv.Atoi(port) @@ -329,12 +331,31 @@ func (t *Listener) readPcap() { for _, d := range devices { go func(device pcap.Interface) { - handle, err := pcap.OpenLive(device.Name, 65536, true, t.messageExpire) + inactive, err := pcap.NewInactiveHandle(device.Name) if err != nil { log.Println("Pcap Error while opening device", device.Name, err) wg.Done() return } + + if t.timestampType != "" { + if tt, terr := pcap.TimestampSourceFromString(t.timestampType); terr != nil { + log.Println("Supported timestamp types: ", inactive.SupportedTimestamps(), device.Name) + } else if terr := inactive.SetTimestampSource(tt); terr != nil { + log.Println("Supported timestamp types: ", inactive.SupportedTimestamps(), device.Name) + } + } + inactive.SetSnapLen(65536) + inactive.SetTimeout(t.messageExpire) + inactive.SetPromisc(true) + + handle, herr := inactive.Activate() + if herr != nil { + log.Println("PCAP Activate error:", herr) + wg.Done() + return + } + defer handle.Close() t.mu.Lock() diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 839a47a..440624c 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -131,12 +131,6 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { } } - if t.IsIncoming { - t.End = time.Now() - } else { - t.End = time.Now().Add(time.Millisecond) - } - if packet.OrigAck != 0 { t.DataAck = packet.OrigAck } @@ -144,6 +138,10 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { if packet.timestamp.Before(t.Start) { t.Start = packet.timestamp } + + if t.End.IsZero() || t.End.Before(packet.timestamp) { + t.End = packet.timestamp + } } t.checkSeqIntegrity() @@ -226,7 +224,7 @@ func (t *TCPMessage) updateHeadersPacket() { return } -// checkIfComplete returns true if all of the packets that compse the message arrived. +// checkIfComplete returns true if all of the packets that compse the message arrived. func (t *TCPMessage) checkIfComplete() { if t.seqMissing || t.headerPacket == -1 { return diff --git a/settings.go b/settings.go index 4815282..3ae0204 100644 --- a/settings.go +++ b/settings.go @@ -54,6 +54,7 @@ type AppSettings struct { inputRAWRealIPHeader string inputRAWExpire time.Duration inputRAWBpfFilter string + inputRAWTimestampType string middleware string @@ -130,6 +131,8 @@ func init() { flag.StringVar(&Settings.inputRAWBpfFilter, "input-raw-bpf-filter", "", "BPF filter to write custom expressions. Can be useful in case of non standard network interfaces like tunneling or SPAN port. Example: --input-raw-bpf-filter 'dst port 80'") + flag.StringVar(&Settings.inputRAWTimestampType, "input-raw-timestamp-type", "", "Possible values: PCAP_TSTAMP_HOST, PCAP_TSTAMP_HOST_LOWPREC, PCAP_TSTAMP_HOST_HIPREC, PCAP_TSTAMP_ADAPTER, PCAP_TSTAMP_ADAPTER_UNSYNCED. This values not supported on all systems, GoReplay will tell you available values of you put wrong one.") + flag.StringVar(&Settings.middleware, "middleware", "", "Used for modifying traffic using external command") // flag.Var(&Settings.inputHTTP, "input-http", "Read requests from HTTP, should be explicitly sent from your application:\n\t# Listen for http on 9000\n\tgor --input-http :9000 --output-http staging.com") @@ -178,9 +181,9 @@ func init() { flag.Var(&Settings.modifierConfig.headerNegativeFilters, "http-disallow-header", "A regexp to match a specific header against. Requests with matching headers will be dropped:\n\t gor --input-raw :8080 --output-http staging.com --http-disallow-header \"User-Agent: Replayed by Gor\"") - flag.Var(&Settings.modifierConfig.headerBasicAuthFilters, "http-basic-auth-filter", "A regexp to match the decoded basic auth string against. Requests with non-matching headers will be dropped:\n\t gor --input-raw :8080 --output-http staging.com --http-basic-auth-filter \"^customer[0-9].*\"") + flag.Var(&Settings.modifierConfig.headerBasicAuthFilters, "http-basic-auth-filter", "A regexp to match the decoded basic auth string against. Requests with non-matching headers will be dropped:\n\t gor --input-raw :8080 --output-http staging.com --http-basic-auth-filter \"^customer[0-9].*\"") - flag.Var(&Settings.modifierConfig.headerHashFilters, "http-header-limiter", "Takes a fraction of requests, consistently taking or rejecting a request based on the FNV32-1A hash of a specific header:\n\t gor --input-raw :8080 --output-http staging.com --http-header-limiter user-id:25%") + flag.Var(&Settings.modifierConfig.headerHashFilters, "http-header-limiter", "Takes a fraction of requests, consistently taking or rejecting a request based on the FNV32-1A hash of a specific header:\n\t gor --input-raw :8080 --output-http staging.com --http-header-limiter user-id:25%") flag.Var(&Settings.modifierConfig.headerHashFilters, "output-http-header-hash-filter", "WARNING: `output-http-header-hash-filter` DEPRECATED, use `--http-header-hash-limiter` instead") From d265594b00ad2dad20cb8c0a75139ebd6a0f6297 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Wed, 7 Feb 2018 20:28:24 +0200 Subject: [PATCH 11/37] Ensure that tcp-output do not loose messages when reconnecting --- Dockerfile.dev | 2 +- output_tcp.go | 8 +++++--- vendor/vendor.json | 3 +-- 3 files changed, 7 insertions(+), 6 deletions(-) diff --git a/Dockerfile.dev b/Dockerfile.dev index fd7850c..da1cfbf 100644 --- a/Dockerfile.dev +++ b/Dockerfile.dev @@ -11,7 +11,7 @@ RUN echo oracle-java7-installer shared/accepted-oracle-license-v1-1 select true RUN apt-get install oracle-java8-installer -y RUN apt-get install flex bison -y -RUN wget http://www.tcpdump.org/release/libpcap-1.7.4.tar.gz && tar xzf libpcap-1.7.4.tar.gz && cd libpcap-1.7.4 && ./configure && make install +RUN wget http://www.tcpdump.org/release/libpcap-1.8.1.tar.gz && tar xzf libpcap-1.8.1.tar.gz && cd libpcap-1.8.1 && ./configure && make install RUN go get github.com/google/gopacket RUN go get -u github.com/golang/lint/golint diff --git a/output_tcp.go b/output_tcp.go index 4c5b9ef..916ba7d 100644 --- a/output_tcp.go +++ b/output_tcp.go @@ -32,7 +32,7 @@ func NewTCPOutput(address string, config *TCPOutputConfig) io.Writer { o.address = address o.config = config - o.buf = make(chan []byte, 100) + o.buf = make(chan []byte, 1000) if Settings.outputTCPStats { o.bufStats = NewGorStat("output_tcp") } @@ -66,11 +66,13 @@ func (o *TCPOutput) worker() { defer conn.Close() for { - conn.Write(<-o.buf) + data := <-o.buf + conn.Write(data) _, err := conn.Write([]byte(payloadSeparator)) if err != nil { - log.Println("Lost connection with aggregator instance, reconnecting") + log.Println("INFO: TCP output connection closed, reconnecting") + o.buf <- data go o.worker() break } diff --git a/vendor/vendor.json b/vendor/vendor.json index 6da1d5c..a6319c5 100644 --- a/vendor/vendor.json +++ b/vendor/vendor.json @@ -46,7 +46,6 @@ }, { "checksumSHA1": "U2Ydh7vEAKlN0Wq22n1JpefF7uY=", - "origin": "github.com/buger/goreplay/vendor/github.com/google/gopacket", "path": "github.com/google/gopacket", "revision": "b09bf408520f7646e29b7033d9adb00ed779a1c4", "revisionTime": "2016-05-12T15:06:07Z" @@ -94,5 +93,5 @@ "revisionTime": "2017-02-01T04:15:14Z" } ], - "rootPath": "github.com/buger/gor" + "rootPath": "github.com/buger/goreplay" } From 1f62555ca860086bbbbcf64ba4017db716a2a9c4 Mon Sep 17 00:00:00 2001 From: Rajat Verma Date: Tue, 27 Feb 2018 06:21:14 -0800 Subject: [PATCH 12/37] Update README.md Fix minor typos --- README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 8cd1283..cfe3db6 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ GoReplay is the simplest and safest way to test your app using real traffic before you put it into production. -As your application grows, the effort required to test it also grows exponentially. GoReplay offers you the simple idea of reusing your existing traffic for testing, which makes it incredibly powerful. Our state of art technique allows to analyze and record your application traffic without affecting it. This eliminates the risks that come with putting a third party component in the critical path. +As your application grows, the effort required to test it also grows exponentially. GoReplay offers you the simple idea of reusing your existing traffic for testing, which makes it incredibly powerful. Our state of art technique allows you to analyze and record your application traffic without affecting it. This eliminates the risks that come with putting a third party component in the critical path. GoReplay increases your confidence in code deployments, configuration changes and infrastructure changes. Did we mention that no coding is required? @@ -68,7 +68,7 @@ If you have a private question feel free to send email to support@gortool.com. * [Granify](http://granify.com) - AI backed SaaS solution that enables online retailers to maximise their sales * And many more! -If you are using Gor we are happy add you to the list and share your story, just write to: hello@goreplay.org +If you are using Gor, we are happy to add you to the list and share your story, just write to: hello@goreplay.org ## Author From a65256d70731b58ba5c42d1f3a5c2f13c6416409 Mon Sep 17 00:00:00 2001 From: Ubuntu Date: Tue, 27 Feb 2018 12:25:29 +0000 Subject: [PATCH 13/37] Fixed kafka client run of brokers and index out of range issues --- emitter.go | 3 ++- input_kafka.go | 4 +++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/emitter.go b/emitter.go index 8e4180c..9a59735 100644 --- a/emitter.go +++ b/emitter.go @@ -61,7 +61,8 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) { if nr > 0 && len(buf) > nr { payload := buf[:nr] meta := payloadMeta(payload) - requestID := string(meta[1]) + // requestID := string(meta[1]) + requestID := string(meta[0]) _maxN := nr if nr > 500 { diff --git a/input_kafka.go b/input_kafka.go index e03fdbb..343cd11 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -3,6 +3,7 @@ package main import ( "encoding/json" "log" + "strings" "github.com/Shopify/sarama" "github.com/Shopify/sarama/mocks" @@ -27,7 +28,8 @@ func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput { con = config.consumer } else { var err error - con, err = sarama.NewConsumer([]string{config.host}, c) + //con, err = sarama.NewConsumer([]string{config.host}, c) + con, err = sarama.NewConsumer(strings.Split(config.host,","), c) if err != nil { log.Fatalln("Failed to start Sarama(Kafka) consumer:", err) From 2db143aaef08bce33ff1531b5bc06de468b75c10 Mon Sep 17 00:00:00 2001 From: Oleg Blokhin Date: Tue, 24 Apr 2018 16:36:36 +0300 Subject: [PATCH 14/37] Update middleware.js fix request filtering in middleware.js --- middleware/middleware.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index 2548c7b..c29b101 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -45,7 +45,7 @@ function init() { if (proxy.ch[chanID]) { proxy.ch[chanID].forEach(function(ch){ let r = ch.cb(msg); - if (r) resp = r; // If one of callback decided not to send response back, do not override it in global callbacks + if (!r) resp = r; // If one of callback decided not to send response back, do not override it in global callbacks }) // Cleanup Individual message channels to avoid memory leaks From cb1ae0e7c0dbad67307ac8e99e988b47a9c3b7d2 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Tue, 24 Apr 2018 19:42:19 +0300 Subject: [PATCH 15/37] Revert "Update middleware.js" This reverts commit 2db143aaef08bce33ff1531b5bc06de468b75c10. --- middleware/middleware.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index c29b101..2548c7b 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -45,7 +45,7 @@ function init() { if (proxy.ch[chanID]) { proxy.ch[chanID].forEach(function(ch){ let r = ch.cb(msg); - if (!r) resp = r; // If one of callback decided not to send response back, do not override it in global callbacks + if (r) resp = r; // If one of callback decided not to send response back, do not override it in global callbacks }) // Cleanup Individual message channels to avoid memory leaks From e9294aeb5f4598db05e02b896c1b8e47cb6b5815 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Thu, 10 May 2018 22:18:47 +0300 Subject: [PATCH 16/37] Fix emitter typo --- emitter.go | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/emitter.go b/emitter.go index 9a59735..8e4180c 100644 --- a/emitter.go +++ b/emitter.go @@ -61,8 +61,7 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) { if nr > 0 && len(buf) > nr { payload := buf[:nr] meta := payloadMeta(payload) - // requestID := string(meta[1]) - requestID := string(meta[0]) + requestID := string(meta[1]) _maxN := nr if nr > 500 { From c9fe4a11ade8c40dd28878827daa9139975ae005 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Thu, 10 May 2018 22:19:08 +0300 Subject: [PATCH 17/37] Fix proto.Body when body < 4 bytes --- proto/proto.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/proto/proto.go b/proto/proto.go index 107de3f..2722527 100644 --- a/proto/proto.go +++ b/proto/proto.go @@ -338,6 +338,10 @@ func DeleteHeader(payload, name []byte) []byte { // Body returns request/response body func Body(payload []byte) []byte { // 4 -> len(EMPTY_LINE) + if len(payload) < 4 { + return []byte{} + } + return payload[MIMEHeadersEndPos(payload):] } From fe6cfd90feb88199120a862046027fc36e4433ea Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Thu, 10 May 2018 22:19:56 +0300 Subject: [PATCH 18/37] Fix responses when header end line separated by multiple packets --- raw_socket_listener/listener.go | 6 ++---- raw_socket_listener/listener_test.go | 24 ++++++++++++------------ raw_socket_listener/tcp_message.go | 19 ++++++++++++++++--- 3 files changed, 30 insertions(+), 19 deletions(-) diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index 47a1aba..ad1e373 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -445,12 +445,12 @@ func (t *Listener) readPcap() { of = 4 case layers.LinkTypeLoop: of = 4 - case layers.LinkTypeRaw: + case layers.LinkTypeRaw, layers.LayerTypeIPv4: of = 0 case layers.LinkTypeLinuxSLL: of = 16 default: - log.Println("Unknown packet layer", packet) + log.Println("Unknown packet layer", decoder, packet) break } @@ -714,8 +714,6 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { } }() - // log.Println("PACKET:", packet, t.seqWithData) - var responseRequest *TCPMessage var message *TCPMessage diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 988cd65..b606540 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -12,7 +12,7 @@ import ( func TestRawListenerInput(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + 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"), time.Now()) @@ -77,7 +77,7 @@ func responsePacket(prev *TCPPacket, payload []byte) *TCPPacket { } func TestSingleAck100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -96,7 +96,7 @@ func TestSingleAck100Continue(t *testing.T) { } func Test100ContinueWithoutWaiting(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -112,7 +112,7 @@ func Test100ContinueWithoutWaiting(t *testing.T) { // Client first sends data without waiting 100-continue, but once response received, generate packets based on Ack payload func Test100ContinueMixed(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 12\r\n\r\n")) @@ -130,7 +130,7 @@ func Test100ContinueMixed(t *testing.T) { } func TestDoubleAck100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -153,7 +153,7 @@ func TestDoubleAck100Continue(t *testing.T) { func TestRawListenerInputResponseByClose(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + 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"), time.Now()) @@ -193,7 +193,7 @@ func TestRawListenerInputResponseByClose(t *testing.T) { func TestRawListenerInputWithoutResponse(t *testing.T) { var req *TCPMessage - listener := NewListener("", "0", EnginePcap, false, 10*time.Millisecond, "") + 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"), time.Now()) @@ -215,7 +215,7 @@ func TestRawListenerInputWithoutResponse(t *testing.T) { func TestRawListenerResponse(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() reqPacket := firstPacket([]byte("GET / HTTP/1.1\r\n\r\n")) @@ -263,7 +263,7 @@ func get100ContinuePackets() (req []*TCPPacket, resp []*TCPPacket) { } func TestShort100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() req, resp := get100ContinuePackets() @@ -275,7 +275,7 @@ func TestShort100Continue(t *testing.T) { // Response comes before Request func Test100ContinueWrongOrder(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() req, resp := get100ContinuePackets() @@ -428,7 +428,7 @@ func permutation(n int, list []*TCPPacket) []*TCPPacket { // Response comes before Request func TestRawListenerChunkedWrongOrder(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n")) @@ -498,7 +498,7 @@ func getMessage() []*TCPPacket { // Response comes before Request func TestRawListenerBench(t *testing.T) { - l := NewListener("", "0", EnginePcap, true, 200*time.Millisecond, "") + l := NewListener("", "0", EnginePcap, true, 200*time.Millisecond, "", "") defer l.Close() // Should re-construct message from all possible combinations diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 440624c..18704bb 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -199,6 +199,7 @@ func (t *TCPMessage) checkSeqIntegrity() { } var bEmptyLine = []byte("\r\n\r\n") +var bBR = []byte("\r\n") var bChunkEnd = []byte("0\r\n\r\n") func (t *TCPMessage) updateHeadersPacket() { @@ -215,9 +216,16 @@ func (t *TCPMessage) updateHeadersPacket() { } for i, p := range t.packets { - if bytes.LastIndex(p.Data, bEmptyLine) != -1 { - t.headerPacket = i - return + if len(p.Data) >= len(bEmptyLine) { + if bytes.LastIndex(p.Data, bEmptyLine) != -1 { + t.headerPacket = i + return + } + } else if bytes.Equal(p.Data, bBR) { + if bytes.LastIndex(t.packets[i-1].Data, bBR) != -1 { + t.headerPacket = i + return + } } } @@ -227,18 +235,23 @@ func (t *TCPMessage) updateHeadersPacket() { // checkIfComplete returns true if all of the packets that compse the message arrived. func (t *TCPMessage) checkIfComplete() { if t.seqMissing || t.headerPacket == -1 { + // log.Println("Seq missing", t.seqMissing, t.packets) return } if t.methodType == httpMethodNotFound { + // log.Println("Method missing", t.methodType, t.packets) return } // Responses can be emitted only if we found request if !t.IsIncoming && t.AssocMessage == nil { + // log.Println("Assoc not found", t) return } + // log.Println("Found?", t) + switch t.bodyType { case httpBodyEmpty: t.complete = true From 0c0ef97b6740b1d9b2857c8dceb0429ea5afbf45 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Sat, 12 May 2018 15:17:26 +0300 Subject: [PATCH 19/37] Add additional tests --- raw_socket_listener/listener_test.go | 41 ++++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index b606540..21f3d33 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -546,3 +546,44 @@ func TestRawListenerBench(t *testing.T) { } } } + +func TestResponseZeroContentLength(t *testing.T) { + var req, resp *TCPMessage + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + defer listener.Close() + + reqPacket := firstPacket([]byte("POST /api/setup/install HTTP/1.1\r\nHost: localhost:22936\r\nUser-Agent: curl/7.57.0\r\nAccept: */*\r\nContent-Length: 0\r\nContent-Type: application/x-www-form-urlencoded\r\n\r\n")) + respPacket := responsePacket(reqPacket, []byte("HTTP/1.1 200\r\nDate: Fri, 11 May 2018 15:09:10 GMT\r\nServer: Kestrel\r\nCache-Control: no-cache\r\nTransfer-Encoding: chunked\r\n\r\n")) + respPacket2 := nextPacket(respPacket, []byte("0\r\n\r\n")) + + // If response packet comes before request + listener.packetsChan <- reqPacket.dump() + listener.packetsChan <- respPacket.dump() + listener.packetsChan <- respPacket2.dump() + + select { + case req = <-listener.messagesChan: + case <-time.After(time.Millisecond): + t.Error("Should return respose immediately") + return + } + + if !req.IsIncoming { + t.Error("Should be request") + } + + select { + case resp = <-listener.messagesChan: + case <-time.After(time.Millisecond): + t.Error("Should return response immediately") + return + } + + if resp.IsIncoming { + t.Error("Should be response") + } + + if !bytes.Equal(resp.UUID(), req.UUID()) { + t.Error("Resp and Req UUID should be equal") + } +} From b1092d156762bd94061b44987ff838be2e3365ee Mon Sep 17 00:00:00 2001 From: Daniel Roop Date: Fri, 18 May 2018 21:39:49 -0400 Subject: [PATCH 20/37] Created helper to return all http headers for node middleware --- middleware/middleware.js | 39 ++++++++++++++++++++++++++++++++++++++- 1 file changed, 38 insertions(+), 1 deletion(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index 2548c7b..475e4fe 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -253,6 +253,21 @@ function setHttpStatus(payload, newStatus) { return setHttpPath(payload, newStatus); } +function httpHeaders(payload) { + var httpHeaderString = payload.slice(0,payload.indexOf("\r\n\r\n") + 4).toString().split("\n").slice(1); + var headers = {}; + + for (var item in httpHeaderString) { + var parts = httpHeaderString[item].split(":"); + + if (parts.length > 1) { + headers[parts[0]] = parts.slice(1).join(":").trim(); + } + } + + return headers; +} + function httpHeader(payload, name) { var currentLine = 0; var i = 0; @@ -394,7 +409,7 @@ module.exports = { // =========== Tests ============== function testRunner(){ - ["init", "parseMessage", "httpMethod", "httpPath", "setHttpHeader", "httpPathParam", "httpHeader", "httpBody", "setHttpBody", "httpBodyParam", "httpCookie", "setHttpCookie"].forEach(function(t){ + ["init", "parseMessage", "httpMethod", "httpPath", "setHttpHeader", "httpPathParam", "httpHeader", "httpBody", "setHttpBody", "httpBodyParam", "httpCookie", "setHttpCookie", "httpHeaders"].forEach(function(t){ console.log(`====== Start ${t} =======`) eval(`TEST_${t}()`) console.log(`====== End ${t} =======`) @@ -668,3 +683,25 @@ function TEST_setHttpCookie() { return fail(`Should add new cookie: ${p}`) } } + +function TEST_httpHeaders() { + const examplePayload = "GET / HTTP/1.1\r\nHost: localhost:3000\r\nUser-Agent: Node\r\nContent-Length:5\r\n\r\nhello"; + + let expectedHeaders = {"Host": "localhost:3000", "User-Agent": "Node", "Content-Length": "5"} + let payload = Buffer.from(examplePayload); + let headers = httpHeaders(payload); + + ["Host", "User-Agent", "Content-Length"].forEach(function(header){ + let actual = headers[header]; + let expected = expectedHeaders[header]; + + if (!actual) { + fail(`${header} Header was not found`); + } + + if (actual != expected) { + fail(`${header} Header not Equal to Expected: ${expected} was ${actual}`); + } + + }) +} From 4e46bb46af23f80ebc05f7ff45b7edb6d6ce40d4 Mon Sep 17 00:00:00 2001 From: Daniel Roop Date: Fri, 18 May 2018 21:59:42 -0400 Subject: [PATCH 21/37] Updated README to include httpHeaders --- middleware/README.md | 1 + 1 file changed, 1 insertion(+) diff --git a/middleware/README.md b/middleware/README.md index b53a0f2..4fb231b 100644 --- a/middleware/README.md +++ b/middleware/README.md @@ -103,6 +103,7 @@ Package expose following functions to process raw HTTP payloads: * `httpPathParam` - get param from URL path: `gor.httpPathParam(req.http, queryParam)` * `setHttpPathParam` - set URL param: `req.http = gor.setHttpPathParam(req.http, queryParam, value)` * `httpStatus` - response status code +* `httpHeaders` - get all headers: `gor.httpHeaders(req.http)` * `httpHeader` - get HTTP header: `gor.httpHeader(req.http, "Content-Length")` * `setHttpHeader` - Set HTTP header, returns modified payload: `req.http = gor.setHttpHeader(req.http, "X-Replayed", "1")` * `httpBody` - get HTTP Body: `gor.httpBody(req.http)` From 7b1544157f5264354575f7c667afea240f6eba3d Mon Sep 17 00:00:00 2001 From: defa sun Date: Mon, 14 May 2018 14:25:00 -0700 Subject: [PATCH 22/37] fix out of range index in headerIndex, github.com/buger/goreplay/issues/578 --- proto/proto.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/proto/proto.go b/proto/proto.go index 2722527..9c5adfe 100644 --- a/proto/proto.go +++ b/proto/proto.go @@ -93,6 +93,11 @@ func headerIndex(payload []byte, name []byte) int { return i - len(name) } + // We are at the end + if i == len(payload) { + return -1 + } + if payload[i] != name[j] { break } From 0fea94a76c26a2350cd50467a1ee33ea2874f96f Mon Sep 17 00:00:00 2001 From: bruce34 Date: Wed, 23 May 2018 18:49:02 +0100 Subject: [PATCH 23/37] Fix 100 continue (#562) * On write error, disconnect to force a new connection * Don't loose bytes checking connection isAlive * Soak up any 1xx return codes, catching up to the real return code --- http_client.go | 34 ++++++++++++++++++++-------------- 1 file changed, 20 insertions(+), 14 deletions(-) diff --git a/http_client.go b/http_client.go index c264e4b..b6c27b6 100644 --- a/http_client.go +++ b/http_client.go @@ -113,16 +113,10 @@ func (c *HTTPClient) Disconnect() { } } -func (c *HTTPClient) isAlive() bool { - one := make([]byte, 1) - +func (c *HTTPClient) isAlive(readBytes *int) (bool) { // Ready 1 byte from socket without timeout to check if it not closed c.conn.SetReadDeadline(time.Now().Add(time.Millisecond)) - _, err := c.conn.Read(one) - - if err == nil { - return true - } + n, err := c.conn.Read(c.respBuf[:1]) if err == io.EOF { Debug("[HTTPClient] connection closed, reconnecting") @@ -133,7 +127,10 @@ func (c *HTTPClient) isAlive() bool { Debug("Detected broken pipe.", err) return false } - + if n != 0 { + *readBytes += n + Debug("[HTTPClient] isAlive readBytes ", *readBytes) + } return true } @@ -153,7 +150,8 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { } }() - if c.conn == nil || !c.isAlive() { + var readBytes, n int + if c.conn == nil || !c.isAlive(&readBytes) { Debug("[HTTPClient] Connecting:", c.baseURL) if err = c.Connect(); err != nil { log.Println("[HTTPClient] Connection error:", err) @@ -181,10 +179,10 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { if _, err = c.conn.Write(data); err != nil { Debug("[HTTPClient] Write error:", err, c.baseURL) response = errorPayload(HTTP_TIMEOUT) + c.Disconnect() return } - var readBytes, n int var currentChunk []byte timeout = time.Now().Add(c.config.Timeout) chunked := false @@ -205,13 +203,21 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { currentContentLength += n } else { // If headers are finished - - if bytes.Contains(c.respBuf[:readBytes], proto.EmptyLine) { + var firstEmptyLine = bytes.Index(c.respBuf[:readBytes], proto.EmptyLine); + if firstEmptyLine != -1 { if bytes.Equal(proto.Header(c.respBuf[:readBytes], []byte("Transfer-Encoding")), []byte("chunked")) { chunked = true } else { status, _ := strconv.Atoi(string(proto.Status(c.respBuf[:readBytes]))) - if (status >= 100 && status < 200) || status == 204 || status == 304 { + // We want to soak up all 100 Continues received to get the real result code + if status >= 100 && status < 200 { + timeout = time.Now().Add(c.config.Timeout) + var deleteLen = firstEmptyLine + len(proto.EmptyLine) + copy(c.respBuf, c.respBuf[deleteLen:readBytes]) + readBytes -= deleteLen + chunks-- + continue + } else if status == 204 || status == 304 { contentLength = 0 break } else { From 18610d3555d38c98f1c54693c09e691752f2eb22 Mon Sep 17 00:00:00 2001 From: Jordan Crawford Date: Thu, 30 Nov 2017 17:07:39 -0600 Subject: [PATCH 24/37] Solve issue 535 missing responses to HEAD requests, and add a test to verify --- raw_socket_listener/listener_test.go | 34 ++++++++++++++++++++++++++++ raw_socket_listener/tcp_message.go | 12 +++++++++- 2 files changed, 45 insertions(+), 1 deletion(-) diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 21f3d33..a9bf502 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -76,6 +76,40 @@ func responsePacket(prev *TCPPacket, payload []byte) *TCPPacket { ) } +func TestHEADRequestNoBody(t *testing.T) { + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + defer listener.Close() + + reqPacket := firstPacket([]byte("HEAD / HTTP/1.1\r\nContent-Length: 0\r\n\r\n")) + respPacket := responsePacket(reqPacket, []byte("HTTP/1.1 200 OK\r\nContent-Length: 100\r\n\r\n")) + + listener.packetsChan <- reqPacket.dump() + listener.packetsChan <- respPacket.dump() + + var req, resp *TCPMessage + select { + case req = <-listener.messagesChan: + case <-time.After( time.Millisecond): + t.Error("Should return request immediately") + return + } + + if !req.IsIncoming { + t.Error("Should be request") + } + + select { + case resp = <-listener.messagesChan: + case <-time.After(20 * time.Millisecond): + t.Error("Should return response immediately") + return + } + + if resp.IsIncoming { + t.Error("Should be response") + } +} + func TestSingleAck100Continue(t *testing.T) { listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") defer listener.Close() diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 18704bb..964770c 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -369,9 +369,19 @@ func (t *TCPMessage) updateBodyType() { case httpMethodNotFound: return case httpMethodKnown: + + if ! t.IsIncoming && + t.AssocMessage != nil && + bytes.IndexByte( t.AssocMessage.Bytes(), ' ') > -1 && + bytes.Equal( []byte("HEAD"), proto.Method(t.AssocMessage.Bytes()) ) { + // Need to check if this is a response to a head request, + // in which case the body has to be empty regardless. + t.bodyType = httpBodyEmpty + return + } + if len(lengthB) > 0 { t.contentLength, _ = strconv.Atoi(string(lengthB)) - if t.contentLength == 0 { t.bodyType = httpBodyEmpty } else { From c5d1112e7ab3147fd985e9cafc855d75bf585755 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Sun, 27 May 2018 12:31:47 +0300 Subject: [PATCH 25/37] Add way to control packet capture buffer size and optimize snaplen MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Added `—input-raw-buffer-size` - Controls size of the OS buffer (in bytes) which holds packets until they dispatched. Default value depends by system: in Linux around 2MB. If you see big package drop, increase this value. Additionally snaplen (max number of bytes being read for each packet) now dynamically set based on interface MTU + max header size. In most situations it should reduce package drop, because each packet will consume less space in buffer. --- elasticsearch_test.go | 2 +- http_client.go | 4 ++-- input_kafka.go | 2 +- input_raw.go | 6 ++++-- input_raw_test.go | 14 ++++++------- limiter.go | 2 +- middleware_test.go | 4 ++-- plugins.go | 2 +- raw_socket_listener/listener.go | 18 +++++++++++++++-- raw_socket_listener/listener_test.go | 30 ++++++++++++++-------------- raw_socket_listener/tcp_message.go | 6 +++--- settings.go | 3 +++ vendor/vendor.json | 6 +++--- 13 files changed, 59 insertions(+), 40 deletions(-) diff --git a/elasticsearch_test.go b/elasticsearch_test.go index 27e624b..7bddcff 100644 --- a/elasticsearch_test.go +++ b/elasticsearch_test.go @@ -28,7 +28,7 @@ func assertExpectedError(returnedError error, t *testing.T) { func assertNoError(returnedError error, t *testing.T) { if nil != returnedError { - t.Errorf("Expected err %s but got %s", nil, returnedError) + t.Errorf("Expected no err but got %s", returnedError) } } diff --git a/http_client.go b/http_client.go index b6c27b6..02285fc 100644 --- a/http_client.go +++ b/http_client.go @@ -113,7 +113,7 @@ func (c *HTTPClient) Disconnect() { } } -func (c *HTTPClient) isAlive(readBytes *int) (bool) { +func (c *HTTPClient) isAlive(readBytes *int) bool { // Ready 1 byte from socket without timeout to check if it not closed c.conn.SetReadDeadline(time.Now().Add(time.Millisecond)) n, err := c.conn.Read(c.respBuf[:1]) @@ -203,7 +203,7 @@ func (c *HTTPClient) Send(data []byte) (response []byte, err error) { currentContentLength += n } else { // If headers are finished - var firstEmptyLine = bytes.Index(c.respBuf[:readBytes], proto.EmptyLine); + var firstEmptyLine = bytes.Index(c.respBuf[:readBytes], proto.EmptyLine) if firstEmptyLine != -1 { if bytes.Equal(proto.Header(c.respBuf[:readBytes], []byte("Transfer-Encoding")), []byte("chunked")) { chunked = true diff --git a/input_kafka.go b/input_kafka.go index 343cd11..90c17cd 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -29,7 +29,7 @@ func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput { } else { var err error //con, err = sarama.NewConsumer([]string{config.host}, c) - con, err = sarama.NewConsumer(strings.Split(config.host,","), c) + con, err = sarama.NewConsumer(strings.Split(config.host, ","), c) if err != nil { log.Fatalln("Failed to start Sarama(Kafka) consumer:", err) diff --git a/input_raw.go b/input_raw.go index 7dddce1..1cd49c6 100644 --- a/input_raw.go +++ b/input_raw.go @@ -21,6 +21,7 @@ type RAWInput struct { listener *raw.Listener bpfFilter string timestampType string + bufferSize int } // Available engines for intercepting traffic @@ -31,7 +32,7 @@ const ( ) // NewRAWInput constructor for RAWInput. Accepts address with port as argument. -func NewRAWInput(address string, engine int, trackResponse bool, expire time.Duration, realIPHeader string, bpfFilter string, timestampType string) (i *RAWInput) { +func NewRAWInput(address string, engine int, trackResponse bool, expire time.Duration, realIPHeader string, bpfFilter string, timestampType string, bufferSize int) (i *RAWInput) { i = new(RAWInput) i.data = make(chan *raw.TCPMessage) i.address = address @@ -42,6 +43,7 @@ func NewRAWInput(address string, engine int, trackResponse bool, expire time.Dur i.quit = make(chan bool) i.trackResponse = trackResponse i.timestampType = timestampType + i.bufferSize = bufferSize i.listen(address) i.listener.IsReady() @@ -79,7 +81,7 @@ func (i *RAWInput) listen(address string) { log.Fatal("input-raw: error while parsing address", err) } - i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.bpfFilter, i.timestampType) + i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.bpfFilter, i.timestampType, i.bufferSize) ch := i.listener.Receiver() diff --git a/input_raw_test.go b/input_raw_test.go index cd0b319..517bdbd 100644 --- a/input_raw_test.go +++ b/input_raw_test.go @@ -44,7 +44,7 @@ func TestRAWInputIPv4(t *testing.T) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "X-Real-IP", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "X-Real-IP", "", "", 0) defer input.Close() output := NewTestOutput(func(data []byte) { @@ -106,7 +106,7 @@ func TestRAWInputNoKeepAlive(t *testing.T) { originAddr := listener.Addr().String() - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() output := NewTestOutput(func(data []byte) { @@ -152,7 +152,7 @@ func TestRAWInputIPv6(t *testing.T) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() output := NewTestOutput(func(data []byte) { @@ -203,7 +203,7 @@ func TestInputRAW100Expect(t *testing.T) { originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "", 0) defer input.Close() // We will use it to get content of raw HTTP request @@ -266,7 +266,7 @@ func TestInputRAWChunkedEncoding(t *testing.T) { })) originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, time.Second, "", "", "", 0) defer input.Close() replay := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -330,7 +330,7 @@ func TestInputRAWLargePayload(t *testing.T) { })) originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() replay := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { @@ -375,7 +375,7 @@ func BenchmarkRAWInput(b *testing.B) { var respCounter, reqCounter int64 - input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(originAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() output := NewTestOutput(func(data []byte) { diff --git a/limiter.go b/limiter.go index cb1a36c..c28ff62 100644 --- a/limiter.go +++ b/limiter.go @@ -96,5 +96,5 @@ func (l *Limiter) Read(data []byte) (n int, err error) { } func (l *Limiter) String() string { - return fmt.Sprintf("Limiting %s to: %d (isPercent: %b)", l.plugin, l.limit, l.isPercent) + return fmt.Sprintf("Limiting %s to: %d (isPercent: %v)", l.plugin, l.limit, l.isPercent) } diff --git a/middleware_test.go b/middleware_test.go index b8d3c9a..28f65dc 100644 --- a/middleware_test.go +++ b/middleware_test.go @@ -118,7 +118,7 @@ func TestEchoMiddleware(t *testing.T) { // Catch traffic from one service fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() // And redirect to another @@ -180,7 +180,7 @@ func TestTokenMiddleware(t *testing.T) { fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) // Catch traffic from one service - input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "") + input := NewRAWInput(fromAddr, EnginePcap, true, testRawExpire, "", "", "", 0) defer input.Close() // And redirect to another diff --git a/plugins.go b/plugins.go index 1cb0637..7cc2f7a 100644 --- a/plugins.go +++ b/plugins.go @@ -106,7 +106,7 @@ func InitPlugins() { } for _, options := range Settings.inputRAW { - registerPlugin(NewRAWInput, options, engine, Settings.inputRAWTrackResponse, Settings.inputRAWExpire, Settings.inputRAWRealIPHeader, Settings.inputRAWBpfFilter, Settings.inputRAWTimestampType) + registerPlugin(NewRAWInput, options, engine, Settings.inputRAWTrackResponse, Settings.inputRAWExpire, Settings.inputRAWRealIPHeader, Settings.inputRAWBpfFilter, Settings.inputRAWTimestampType, Settings.inputRawBufferSize) } for _, options := range Settings.inputTCP { diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index ad1e373..f4d7246 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -75,6 +75,8 @@ type Listener struct { bpfFilter string timestampType string + bufferSize int + conn net.PacketConn pcapHandles []*pcap.Handle @@ -96,7 +98,7 @@ const ( ) // NewListener creates and initializes new Listener object -func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration, bpfFilter string, timestampType string) (l *Listener) { +func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration, bpfFilter string, timestampType string, bufferSize int) (l *Listener) { l = &Listener{} l.packetsChan = make(chan *packet, 10000) @@ -112,6 +114,7 @@ func NewListener(addr string, port string, engine int, trackResponse bool, expir l.trackResponse = trackResponse l.bpfFilter = bpfFilter l.timestampType = timestampType + l.bufferSize = bufferSize l.addr = addr _port, _ := strconv.Atoi(port) @@ -345,10 +348,21 @@ func (t *Listener) readPcap() { log.Println("Supported timestamp types: ", inactive.SupportedTimestamps(), device.Name) } } - inactive.SetSnapLen(65536) + + if it, err := net.InterfaceByName(device.Name); err == nil { + // Auto-guess max length of packet to capture + inactive.SetSnapLen(it.MTU + 68*2) + } else { + inactive.SetSnapLen(65536) + } + inactive.SetTimeout(t.messageExpire) inactive.SetPromisc(true) + if t.bufferSize > 0 { + inactive.SetBufferSize(t.bufferSize) + } + handle, herr := inactive.Activate() if herr != nil { log.Println("PCAP Activate error:", herr) diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index a9bf502..2a6d9c0 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -12,7 +12,7 @@ import ( func TestRawListenerInput(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) @@ -77,7 +77,7 @@ func responsePacket(prev *TCPPacket, payload []byte) *TCPPacket { } func TestHEADRequestNoBody(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := firstPacket([]byte("HEAD / HTTP/1.1\r\nContent-Length: 0\r\n\r\n")) @@ -89,7 +89,7 @@ func TestHEADRequestNoBody(t *testing.T) { var req, resp *TCPMessage select { case req = <-listener.messagesChan: - case <-time.After( time.Millisecond): + case <-time.After(time.Millisecond): t.Error("Should return request immediately") return } @@ -111,7 +111,7 @@ func TestHEADRequestNoBody(t *testing.T) { } func TestSingleAck100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -130,7 +130,7 @@ func TestSingleAck100Continue(t *testing.T) { } func Test100ContinueWithoutWaiting(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -146,7 +146,7 @@ func Test100ContinueWithoutWaiting(t *testing.T) { // Client first sends data without waiting 100-continue, but once response received, generate packets based on Ack payload func Test100ContinueMixed(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() req1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 12\r\n\r\n")) @@ -164,7 +164,7 @@ func Test100ContinueMixed(t *testing.T) { } func TestDoubleAck100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nExpect: 100-continue\r\nContent-Length: 4\r\n\r\n")) @@ -187,7 +187,7 @@ func TestDoubleAck100Continue(t *testing.T) { func TestRawListenerInputResponseByClose(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) @@ -227,7 +227,7 @@ func TestRawListenerInputResponseByClose(t *testing.T) { func TestRawListenerInputWithoutResponse(t *testing.T) { var req *TCPMessage - listener := NewListener("", "0", EnginePcap, false, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, false, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) @@ -249,7 +249,7 @@ func TestRawListenerInputWithoutResponse(t *testing.T) { func TestRawListenerResponse(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := firstPacket([]byte("GET / HTTP/1.1\r\n\r\n")) @@ -297,7 +297,7 @@ func get100ContinuePackets() (req []*TCPPacket, resp []*TCPPacket) { } func TestShort100Continue(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() req, resp := get100ContinuePackets() @@ -309,7 +309,7 @@ func TestShort100Continue(t *testing.T) { // Response comes before Request func Test100ContinueWrongOrder(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() req, resp := get100ContinuePackets() @@ -462,7 +462,7 @@ func permutation(n int, list []*TCPPacket) []*TCPPacket { // Response comes before Request func TestRawListenerChunkedWrongOrder(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n")) @@ -532,7 +532,7 @@ func getMessage() []*TCPPacket { // Response comes before Request func TestRawListenerBench(t *testing.T) { - l := NewListener("", "0", EnginePcap, true, 200*time.Millisecond, "", "") + l := NewListener("", "0", EnginePcap, true, 200*time.Millisecond, "", "", 0) defer l.Close() // Should re-construct message from all possible combinations @@ -583,7 +583,7 @@ func TestRawListenerBench(t *testing.T) { func TestResponseZeroContentLength(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "") + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0) defer listener.Close() reqPacket := firstPacket([]byte("POST /api/setup/install HTTP/1.1\r\nHost: localhost:22936\r\nUser-Agent: curl/7.57.0\r\nAccept: */*\r\nContent-Length: 0\r\nContent-Type: application/x-www-form-urlencoded\r\n\r\n")) diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 964770c..cb15b9b 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -370,10 +370,10 @@ func (t *TCPMessage) updateBodyType() { return case httpMethodKnown: - if ! t.IsIncoming && + if !t.IsIncoming && t.AssocMessage != nil && - bytes.IndexByte( t.AssocMessage.Bytes(), ' ') > -1 && - bytes.Equal( []byte("HEAD"), proto.Method(t.AssocMessage.Bytes()) ) { + bytes.IndexByte(t.AssocMessage.Bytes(), ' ') > -1 && + bytes.Equal([]byte("HEAD"), proto.Method(t.AssocMessage.Bytes())) { // Need to check if this is a response to a head request, // in which case the body has to be empty regardless. t.bodyType = httpBodyEmpty diff --git a/settings.go b/settings.go index 3ae0204..edf7ba6 100644 --- a/settings.go +++ b/settings.go @@ -55,6 +55,7 @@ type AppSettings struct { inputRAWExpire time.Duration inputRAWBpfFilter string inputRAWTimestampType string + inputRawBufferSize int middleware string @@ -133,6 +134,8 @@ func init() { flag.StringVar(&Settings.inputRAWTimestampType, "input-raw-timestamp-type", "", "Possible values: PCAP_TSTAMP_HOST, PCAP_TSTAMP_HOST_LOWPREC, PCAP_TSTAMP_HOST_HIPREC, PCAP_TSTAMP_ADAPTER, PCAP_TSTAMP_ADAPTER_UNSYNCED. This values not supported on all systems, GoReplay will tell you available values of you put wrong one.") + flag.IntVar(&Settings.inputRawBufferSize, "input-raw-buffer-size", 0, "Controls size of the OS buffer (in bytes) which holds packets until they dispatched. Default value depends by system: in Linux around 2MB. If you see big package drop, increase this value.") + flag.StringVar(&Settings.middleware, "middleware", "", "Used for modifying traffic using external command") // flag.Var(&Settings.inputHTTP, "input-http", "Read requests from HTTP, should be explicitly sent from your application:\n\t# Listen for http on 9000\n\tgor --input-http :9000 --output-http staging.com") diff --git a/vendor/vendor.json b/vendor/vendor.json index a6319c5..5db34eb 100644 --- a/vendor/vendor.json +++ b/vendor/vendor.json @@ -45,10 +45,10 @@ "revisionTime": "2016-05-29T05:00:41Z" }, { - "checksumSHA1": "U2Ydh7vEAKlN0Wq22n1JpefF7uY=", + "checksumSHA1": "WT6lYgJhoWbXLpnFOxPISxrL2/o=", "path": "github.com/google/gopacket", - "revision": "b09bf408520f7646e29b7033d9adb00ed779a1c4", - "revisionTime": "2016-05-12T15:06:07Z" + "revision": "60ab61cd59496fcfa4d208b265ba79b1e37c1476", + "revisionTime": "2018-05-13T17:29:36Z" }, { "checksumSHA1": "BM6ZlNJmtKy3GBoWwg2X55gnZ4A=", From edae33717e76f2321e2cbae6c63a1fcc765809dd Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Mon, 28 May 2018 10:05:07 +0300 Subject: [PATCH 26/37] Improve start/end time calculation --- raw_socket_listener/listener_test.go | 4 ++-- raw_socket_listener/tcp_message.go | 5 ++--- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 2a6d9c0..e29e25f 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -62,7 +62,7 @@ func nextPacket(prev *TCPPacket, payload []byte) *TCPPacket { prev.Ack, prev.Seq+uint32(len(prev.Data)), payload, - time.Now(), + prev.timestamp.Add(time.Millisecond), ) } @@ -72,7 +72,7 @@ func responsePacket(prev *TCPPacket, payload []byte) *TCPPacket { prev.Seq+uint32(len(prev.Data)), prev.Ack, payload, - time.Now(), + prev.timestamp.Add(time.Millisecond), ) } diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index cb15b9b..27c65dd 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -130,16 +130,15 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { } } } - if packet.OrigAck != 0 { t.DataAck = packet.OrigAck } - if packet.timestamp.Before(t.Start) { + if packet.timestamp.Before(t.Start) || t.Start.IsZero() { t.Start = packet.timestamp } - if t.End.IsZero() || t.End.Before(packet.timestamp) { + if packet.timestamp.After(t.End) || t.End.IsZero() { t.End = packet.timestamp } } From 896dbe312ae2bb339add05237a3b35b870074146 Mon Sep 17 00:00:00 2001 From: Chintan Patel Date: Thu, 7 Jun 2018 19:27:34 -0400 Subject: [PATCH 27/37] Exporting the `httpHeaders` method This method export was missing from last commit (when `httpHeaders` was created) --- middleware/middleware.js | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index 475e4fe..a6feb81 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -402,7 +402,8 @@ module.exports = { httpCookie: httpCookie, setHttpCookie: setHttpCookie, test: testRunner, - benchmark: testBenchmark + benchmark: testBenchmark, + httpHeaders: httpHeaders } From 9eb6ab262b83d29f87ca0711894be6a2d50a4627 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Fri, 22 Jun 2018 22:11:29 +0500 Subject: [PATCH 28/37] Add "deleteHttpHeader" --- middleware/middleware.js | 25 ++++++++++++++++++++++++- 1 file changed, 24 insertions(+), 1 deletion(-) diff --git a/middleware/middleware.js b/middleware/middleware.js index a6feb81..7964a15 100755 --- a/middleware/middleware.js +++ b/middleware/middleware.js @@ -322,6 +322,16 @@ function setHttpHeader(payload, name, value) { } } +function deleteHttpHeader(payload, name) { + let header = httpHeader(payload, name); + + if (header) { + return Buffer.concat([payload.slice(0, header.start), payload.slice(header.end+1, payload.length)]) + } + + return payload +} + function httpBody(payload) { return payload.slice(payload.indexOf("\r\n\r\n") + 4, payload.length); } @@ -395,6 +405,7 @@ module.exports = { setHttpStatus: setHttpStatus, httpHeader: httpHeader, setHttpHeader: setHttpHeader, + deleteHttpHeader: deleteHttpHeader, httpBody: httpBody, setHttpBody: setHttpBody, httpBodyParam: httpBodyParam, @@ -410,7 +421,7 @@ module.exports = { // =========== Tests ============== function testRunner(){ - ["init", "parseMessage", "httpMethod", "httpPath", "setHttpHeader", "httpPathParam", "httpHeader", "httpBody", "setHttpBody", "httpBodyParam", "httpCookie", "setHttpCookie", "httpHeaders"].forEach(function(t){ + ["init", "parseMessage", "httpMethod", "httpPath", "setHttpHeader", "deleteHttpHeader", "httpPathParam", "httpHeader", "httpBody", "setHttpBody", "httpBodyParam", "httpCookie", "setHttpCookie", "httpHeaders"].forEach(function(t){ console.log(`====== Start ${t} =======`) eval(`TEST_${t}()`) console.log(`====== End ${t} =======`) @@ -642,6 +653,18 @@ function TEST_setHttpHeader() { } } +function TEST_deleteHttpHeader() { + const examplePayload = "GET / HTTP/1.1\r\nUser-Agent: Node\r\nContent-Length: 5\r\n\r\nhello"; + + // Adding new header + let expected = `GET / HTTP/1.1\r\nContent-Length: 5\r\n\r\nhello`; + let p = Buffer.from(examplePayload); + p = deleteHttpHeader(p, "User-Agent", "test"); + if (p != expected) { + console.error(`setHeader failed, expected delete header 'User-Agent' header: ${p}`) + } +} + function TEST_httpBody() { const examplePayload = "GET / HTTP/1.1\r\nUser-Agent: Node\r\nContent-Length: 5\r\n\r\nhello"; let body = httpBody(Buffer.from(examplePayload)); From 36a8845ffe8abfbcdab0c28f909ef99780f4f895 Mon Sep 17 00:00:00 2001 From: Alex Zvorygin Date: Wed, 27 Jun 2018 15:10:41 -0400 Subject: [PATCH 29/37] Fix typo --- examples/middleware/echo.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/examples/middleware/echo.rb b/examples/middleware/echo.rb index 75311a8..12bc03c 100755 --- a/examples/middleware/echo.rb +++ b/examples/middleware/echo.rb @@ -6,7 +6,7 @@ while data = STDIN.gets # continiously read line from STDIN decoded = [data].pack("H*") # decode base64 encoded request - # dedoded value is raw HTTP payload, example: + # decoded value is raw HTTP payload, example: # # POST /post HTTP/1.1 # Content-Length: 7 From 85a7bc1e9a74ad5e0d628b5060290668b0cf06de Mon Sep 17 00:00:00 2001 From: Samuel Gulliksson Date: Wed, 27 Jun 2018 16:03:17 +0200 Subject: [PATCH 30/37] Fix typo in Python middleware example. --- examples/middleware/echo.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/examples/middleware/echo.py b/examples/middleware/echo.py index f8b7c99..a38d1af 100644 --- a/examples/middleware/echo.py +++ b/examples/middleware/echo.py @@ -55,7 +55,7 @@ def process_stdin(): request_type_id = int(raw_metadata.split(b' ')[0]) log('Request type: {}'.format({ 1: 'Request', - 2: 'Original Request', + 2: 'Original Response', 3: 'Replayed Response' }[request_type_id])) log('===================================') From 88972ff301ee860708e190880d15123767cb77df Mon Sep 17 00:00:00 2001 From: Alex Zvorygin Date: Tue, 3 Jul 2018 09:19:46 -0400 Subject: [PATCH 31/37] fix typo in ruby middleware example continiously -> continuously --- examples/middleware/echo.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/examples/middleware/echo.rb b/examples/middleware/echo.rb index 12bc03c..7afe673 100755 --- a/examples/middleware/echo.rb +++ b/examples/middleware/echo.rb @@ -1,6 +1,6 @@ #!/usr/bin/env ruby # encoding: utf-8 -while data = STDIN.gets # continiously read line from STDIN +while data = STDIN.gets # continuously read line from STDIN next unless data data = data.chomp # remove end of line symbol From 65fd1983b140a2fcd6c4cc71a2c44582d0879275 Mon Sep 17 00:00:00 2001 From: Alex Zvorygin Date: Mon, 9 Jul 2018 15:48:47 -0400 Subject: [PATCH 32/37] Fix broken link to goreplay.com -> goreplay.org --- docs/Distributed-configuration.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/Distributed-configuration.md b/docs/Distributed-configuration.md index 2795bf8..e08cc28 100644 --- a/docs/Distributed-configuration.md +++ b/docs/Distributed-configuration.md @@ -13,7 +13,7 @@ If you have multiple replay machines you can split traffic among them using `--s gor --input-raw :80 --split-output --output-tcp replay1.local:28020 --output-tcp replay2.local:28020 ``` -[GoReplay PRO](https://goreplay.com/pro.html) support accurate recording and replaying of tcp sessions, and when `--recognize-tcp-sessions` option is passed, instead of round-robin it will use a smarter algorithm which ensures that same sessions will be sent to the same replay instance. +[GoReplay PRO](https://goreplay.org/pro.html) support accurate recording and replaying of tcp sessions, and when `--recognize-tcp-sessions` option is passed, instead of round-robin it will use a smarter algorithm which ensures that same sessions will be sent to the same replay instance. In case if you are planning a large load testing, you may consider use separate master instance which will control Gor slaves which actually replay traffic. For example: From 0bc3fba5d021214bfe40812fd8af696a4bc92b1e Mon Sep 17 00:00:00 2001 From: Scott Albertson Date: Tue, 31 Jul 2018 11:20:42 -0700 Subject: [PATCH 33/37] Add a "Reviewed by Hound" badge --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index cfe3db6..227e8ff 100644 --- a/README.md +++ b/README.md @@ -1,4 +1,4 @@ -[![GitHub release](https://img.shields.io/github/release/buger/gor.svg?maxAge=3600)](https://github.com/buger/goreplay/releases) [![codebeat](https://codebeat.co/badges/6427d589-a78e-416c-a546-d299b4089893)](https://codebeat.co/projects/github-com-buger-gor) [![Go Report Card](https://goreportcard.com/badge/github.com/buger/gor)](https://goreportcard.com/report/github.com/buger/gor) [![Join the chat at https://gitter.im/buger/gor](https://badges.gitter.im/buger/gor.svg)](https://gitter.im/buger/gor?utm_source=badge&utm_medium=badge&utm_campaign=pr-badge&utm_content=badge) +[![GitHub release](https://img.shields.io/github/release/buger/gor.svg?maxAge=3600)](https://github.com/buger/goreplay/releases) [![codebeat](https://codebeat.co/badges/6427d589-a78e-416c-a546-d299b4089893)](https://codebeat.co/projects/github-com-buger-gor) [![Go Report Card](https://goreportcard.com/badge/github.com/buger/gor)](https://goreportcard.com/report/github.com/buger/gor) [![Join the chat at https://gitter.im/buger/gor](https://badges.gitter.im/buger/gor.svg)](https://gitter.im/buger/gor?utm_source=badge&utm_medium=badge&utm_campaign=pr-badge&utm_content=badge) [![Reviewed by Hound](https://img.shields.io/badge/Reviewed_by-Hound-8E64B0.svg)](https://houndci.com) ![Go Replay](http://i.imgur.com/ZG2ki5n.png) From c93ff3598c55bebbab2f8691a03f2b6aecd52fdd Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Thu, 2 Aug 2018 21:36:47 +0500 Subject: [PATCH 34/37] Fix prettifier buffer overflow --- middleware.go | 42 +++++++++++++++++++++++++++--------------- 1 file changed, 27 insertions(+), 15 deletions(-) diff --git a/middleware.go b/middleware.go index 0b720e0..c90270e 100644 --- a/middleware.go +++ b/middleware.go @@ -62,28 +62,40 @@ func (m *Middleware) ReadFrom(plugin io.Reader) { func (m *Middleware) copy(to io.Writer, from io.Reader) { buf := make([]byte, 5*1024*1024) - dst := make([]byte, len(buf)*3) + dst := make([]byte, len(buf)*4) for { nr, _ := from.Read(buf) - if nr > 0 && len(buf) > nr { - payload := buf[0:nr] - - if Settings.prettifyHTTP { - payload = prettifyHTTP(payload) - nr = len(payload) + if nr = 0 || nr > len(buf) { + continue + } + + payload := buf[0:nr] + + if Settings.prettifyHTTP { + payload = prettifyHTTP(payload) + nr = len(payload) + + if nr*2 > len(dst) { + continue } + } + - hex.Encode(dst, payload) - dst[nr*2] = '\n' + if Settings.prettifyHTTP { + payload = prettifyHTTP(payload) + nr = len(payload) + } - m.mu.Lock() - to.Write(dst[0 : nr*2+1]) - m.mu.Unlock() + hex.Encode(dst, payload) + dst[nr*2] = '\n' - if Settings.debug { - Debug("[MIDDLEWARE-MASTER] Sending:", string(buf[0:nr]), "From:", from) - } + m.mu.Lock() + to.Write(dst[0 : nr*2+1]) + m.mu.Unlock() + + if Settings.debug { + Debug("[MIDDLEWARE-MASTER] Sending:", string(buf[0:nr]), "From:", from) } } } From 68082bfb1183dd3d50f5743f41a8a3a444c6adb0 Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Thu, 2 Aug 2018 21:41:36 +0500 Subject: [PATCH 35/37] Fix typo --- middleware.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/middleware.go b/middleware.go index c90270e..0fa0031 100644 --- a/middleware.go +++ b/middleware.go @@ -66,21 +66,20 @@ func (m *Middleware) copy(to io.Writer, from io.Reader) { for { nr, _ := from.Read(buf) - if nr = 0 || nr > len(buf) { + if nr == 0 || nr > len(buf) { continue } - + payload := buf[0:nr] - + if Settings.prettifyHTTP { payload = prettifyHTTP(payload) nr = len(payload) - + if nr*2 > len(dst) { continue } } - if Settings.prettifyHTTP { payload = prettifyHTTP(payload) From f85bb293270d0e049bd6029e36d9aef5ee4fc5aa Mon Sep 17 00:00:00 2001 From: chainhelen Date: Mon, 6 Aug 2018 02:42:16 +0800 Subject: [PATCH 36/37] remove unused "copy" --- 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 f4d7246..e950e74 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -682,12 +682,6 @@ func (t *Listener) readRAWSocket() { } 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, From 9fd7a46e7ecd95543a9ad50af5b8417c210f2808 Mon Sep 17 00:00:00 2001 From: myzhan Date: Thu, 7 Jun 2018 18:38:43 +0800 Subject: [PATCH 37/37] FIX: miscalculation of the mean value --- gor_stat.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gor_stat.go b/gor_stat.go index 205640e..2273a5c 100644 --- a/gor_stat.go +++ b/gor_stat.go @@ -40,7 +40,7 @@ func (s *GorStat) Write(latest int) { s.max = latest } if latest != 0 { - s.mean = (s.mean + latest) / 2 + s.mean = ((s.mean * s.count) + latest) / (s.count + 1) } s.latest = latest s.count = s.count + 1