mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Fix memory leak and improve tests
This commit is contained in:
@@ -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:
|
||||
|
||||
+8
-1
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user