From 6a3b791f7b363033a143fec4dd3df8f870807e5c Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Mon, 9 May 2016 18:00:12 +0500 Subject: [PATCH] Fix memory leak and improve tests --- Makefile | 6 +- input_raw.go | 9 ++- input_raw_test.go | 7 -- raw_socket_listener/listener.go | 22 +++++- raw_socket_listener/listener_test.go | 112 ++++++++++++++++++++++++++- 5 files changed, 142 insertions(+), 14 deletions(-) diff --git a/Makefile b/Makefile index dac397c..420fbe3 100644 --- a/Makefile +++ b/Makefile @@ -2,6 +2,7 @@ SOURCE = emitter.go gor.go gor_stat.go input_dummy.go input_file.go input_raw.go SOURCE_PATH = /go/src/github.com/buger/gor/ RUN = docker run -v `pwd`:$(SOURCE_PATH) -p 0.0.0.0:8000:8000 -t -i gor BENCHMARK = BenchmarkRAWInput +TEST = TestRawListenerBench release: release-x64 @@ -47,9 +48,8 @@ bench: $(RUN) go test -v -run NOT_EXISTING -bench $(BENCHMARK) -benchtime 5s profile_test: - $(RUN) go test $(LDFLAGS) -run NOT_EXISTING -test.benchmem -bench $(BENCHMARK) ./. $(ARGS) -benchtime 5s -memprofile mem.mprof -v - $(RUN) go test $(LDFLAGS) -run NOT_EXISTING -test.benchmem -bench $(BENCHMARK) ./. $(ARGS) -benchtime 5s -cpuprofile cpu.out -v - $(RUN) go test $(LDFLAGS) -run NOT_EXISTING -test.benchmem -bench $(BENCHMARK) ./. $(ARGS) -c + $(RUN) go test $(LDFLAGS) -run $(TEST) ./raw_socket_listener/. $(ARGS) -memprofile mem.mprof -cpuprofile cpu.out + $(RUN) go test $(LDFLAGS) -run $(TEST) ./raw_socket_listener/. $(ARGS) -c # Used mainly for debugging, because docker container do not have access to parent machine ports run: diff --git a/input_raw.go b/input_raw.go index 60a495b..e98d94c 100644 --- a/input_raw.go +++ b/input_raw.go @@ -35,6 +35,11 @@ func NewRAWInput(address string, engine int, expire time.Duration) (i *RAWInput) go i.listen(address) + for i.listener == nil { + time.Sleep(time.Millisecond) + } + i.listener.IsReady() + return } @@ -69,6 +74,8 @@ func (i *RAWInput) listen(address string) { i.listener = raw.NewListener(host, port, i.engine, i.expire) + ch := i.listener.Receiver() + for { select { case <-i.quit: @@ -77,7 +84,7 @@ func (i *RAWInput) listen(address string) { } // Receiving TCPMessage object - m := i.listener.Receive() + m := <- ch i.data <- m } diff --git a/input_raw_test.go b/input_raw_test.go index 9630921..2fb7bf0 100644 --- a/input_raw_test.go +++ b/input_raw_test.go @@ -54,7 +54,6 @@ func TestRAWInput(t *testing.T) { client := NewHTTPClient(origin.URL, &HTTPClientConfig{}) go Start(quit) - time.Sleep(100 * time.Millisecond) for i := 0; i < 100; i++ { // request + response @@ -118,7 +117,6 @@ func TestInputRAW100Expect(t *testing.T) { Plugins.Outputs = []io.Writer{testOutput, httpOutput} go Start(quit) - time.Sleep(100 * time.Millisecond) // Origin + Response/Request Test Output + Request Http Output wg.Add(4) @@ -169,7 +167,6 @@ func TestInputRAWChunkedEncoding(t *testing.T) { Plugins.Outputs = []io.Writer{httpOutput} go Start(quit) - time.Sleep(100 * time.Millisecond) wg.Add(2) @@ -237,8 +234,6 @@ func TestInputRAWLargePayload(t *testing.T) { go Start(quit) - time.Sleep(100 * time.Millisecond) - wg.Add(2) curl := exec.Command("curl", "http://"+originAddr, "--header", "Transfer-Encoding: chunked", "--header", "Expect:", "--data-binary", "@/tmp/large") err = curl.Run() @@ -275,8 +270,6 @@ func BenchmarkRAWInput(b *testing.B) { Plugins.Inputs = []io.Reader{input} Plugins.Outputs = []io.Writer{output} - time.Sleep(time.Millisecond) - go Start(quit) emitted := 0 diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index b579bb9..1327ae2 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -63,6 +63,7 @@ type Listener struct { conn net.PacketConn quit chan bool + readyCh chan bool } type request struct { @@ -84,6 +85,7 @@ func NewListener(addr string, port string, engine int, expire time.Duration) (l l.packetsChan = make(chan *TCPPacket, 10000) l.messagesChan = make(chan *TCPMessage, 10000) l.quit = make(chan bool) + l.readyCh = make(chan bool, 1) l.messages = make(map[string]*TCPMessage) l.ackAliases = make(map[uint32]uint32) @@ -165,6 +167,7 @@ func (t *Listener) dispatchMessage(message *TCPMessage) { delete(t.ackAliases, message.Ack) delete(t.messages, message.ID) + delete(t.respAliases, message.ResponseAck) // log.Println("Dispatching, message", message.Seq, message.Ack, string(message.Bytes())) @@ -184,6 +187,8 @@ func (t *Listener) dispatchMessage(message *TCPMessage) { } } } + + } else { if message.RequestAck == 0 { if responseRequest, ok := t.respAliases[message.Ack]; ok { @@ -265,6 +270,8 @@ func (t *Listener) readPcap() { source.Lazy = true source.NoCopy = true + t.readyCh <- true + // log.Println(handle.Stats()) for { @@ -320,6 +327,8 @@ func (t *Listener) readRAWSocket() { buf := make([]byte, 64*1024) // 64kb + t.readyCh <- true + for { // Note: ReadFrom receive messages without IP header n, addr, err := t.conn.ReadFrom(buf) @@ -500,9 +509,18 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { } } +func (t *Listener) IsReady() bool { + select { + case <- t.readyCh: + return true + case <- time.After(5 * time.Second): + return false + } +} + // Receive TCP messages from the listener channel -func (t *Listener) Receive() *TCPMessage { - return <-t.messagesChan +func (t *Listener) Receiver() chan *TCPMessage { + return t.messagesChan } func (t *Listener) Close() { diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index e52796b..e317423 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -2,9 +2,11 @@ package rawSocket import ( "bytes" - _ "log" + "log" "testing" "time" + "math/rand" + "sync/atomic" ) func TestRawListenerInput(t *testing.T) { @@ -219,8 +221,16 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket select { case r = <-listener.messagesChan: if r.IsIncoming { + if req != nil { + t.Error("Request already received", r) + return + } req = r } else { + if resp != nil { + t.Error("Response already received", r) + return + } resp = r } break @@ -289,3 +299,103 @@ func TestRawListenerChunkedWrongOrder(t *testing.T) { testChunkedSequence(t, listener, packets...) } } + + +func chunkedPostMessage() []*TCPPacket { + ack := uint32(rand.Int63()) + seq := uint32(rand.Int63()) + + reqPacket1 := buildPacket(true, ack, seq, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n")) + // Packet with data have different Seq + reqPacket2 := buildPacket(true, ack, seq + 47, []byte("1\r\na\r\n")) + reqPacket3 := buildPacket(true, ack, reqPacket2.Seq+5, []byte("1\r\nb\r\n")) + reqPacket4 := buildPacket(true, ack, reqPacket3.Seq+5, []byte("0\r\n\r\n")) + + respPacket := buildPacket(false, reqPacket4.Seq+5 /* len of data */, ack, []byte("HTTP/1.1 200 OK\r\n")) + + return []*TCPPacket{ + reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket, + } +} + +func postMessage() []*TCPPacket { + ack := uint32(rand.Int63()) + seq2 := uint32(rand.Int63()) + seq := uint32(rand.Int63()) + + c := 10000 + data := make([]byte, c) + rand.Read(data) + + head := []byte("POST / HTTP/1.1\r\nContent-Length: 9958\r\n\r\n") + for i, _ := range head { + data[i] = head[i] + } + + return []*TCPPacket{ + buildPacket(true, ack, seq, data), + buildPacket(false, seq + uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n")), + } +} + +func getMessage() []*TCPPacket { + ack := uint32(rand.Int63()) + seq2 := uint32(rand.Int63()) + seq := uint32(rand.Int63()) + + return []*TCPPacket{ + buildPacket(true, ack, seq, []byte("GET / HTTP/1.1\r\n\r\n")), + buildPacket(false, seq + 18, seq2, []byte("HTTP/1.1 200 OK\r\n")), + } +} + +// Response comes before Request +func TestRawListenerBench(t *testing.T) { + l := NewListener("", "0", EnginePcap, 200*time.Millisecond) + defer l.Close() + + // Should re-construct message from all possible combinations + for i := 0; i < 1000; i++ { + go func(){ + for j := 0; j < 100; j++ { + var packets []*TCPPacket + + if j % 5 == 0 { + packets = chunkedPostMessage() + } else if j % 3 == 0 { + packets = postMessage() + } else { + packets = getMessage() + } + + for _, p := range packets { + // Randomly drop packets + if (i + j) % 5 == 0 { + if rand.Int63() % 3 == 0 { + continue + } + } + + l.packetsChan <- p + time.Sleep(time.Millisecond) + } + + time.Sleep(5 * time.Millisecond) + } + }() + } + + ch := l.Receiver() + + var count int32 + + for { + select { + case <- ch: + atomic.AddInt32(&count, 1) + case <-time.After(2000 * time.Millisecond): + log.Println("Emitted 200000 messages, captured: ", count, len(l.ackAliases), len(l.seqWithData), len(l.respAliases), len(l.respWithoutReq), len(l.packetsChan)) + return + } + } +}