From 2c34f291fb450ff24173f9fb7ab76e53f18f6fb2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Fri, 9 Aug 2013 21:07:23 +0000 Subject: [PATCH 01/31] ignore vim files --- .gitignore | 1 + 1 file changed, 1 insertion(+) create mode 100644 .gitignore diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..1377554 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +*.swp From 719fd471f36227d52deff54d53f67d87b8258c65 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 19 Aug 2013 18:53:28 +0200 Subject: [PATCH 02/31] failing integration spec --- integration_test.go | 84 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 84 insertions(+) diff --git a/integration_test.go b/integration_test.go index ab445ef..9f48017 100644 --- a/integration_test.go +++ b/integration_test.go @@ -189,3 +189,87 @@ func TestListenerRateLimit(t *testing.T) { t.Error("It should forward only 3 requests with rate-limiting", processed) } } + +func (e *Env) startWithFile() (p int) { + p = 50000 + envs*10 + + go e.startHTTP(p, http.HandlerFunc(e.ListenHandler)) + go e.startHTTP(p+2, http.HandlerFunc(e.ReplayHandler)) + go e.startFileUsingListener(p, p+1) + go e.startFileUsingReplay(p+1, p+2) + + // Time to start http and gor instances + time.Sleep(time.Millisecond * 100) + + envs++ + + return +} + +func (e *Env) startFileUsingListener(port int, replayPort int) { + listener.Settings.Verbose = e.Verbose + listener.Settings.Address = "127.0.0.1" + listener.Settings.FileToReplyPath = "integration_request.gor" + listener.Settings.Port = port + + if e.ListenerLimit != 0 { + listener.Settings.ReplayAddress += "|" + strconv.Itoa(e.ListenerLimit) + } + + listener.Run() +} + +func (e *Env) startFileUsingReplay(port int, forwardPort int) { + replay.Settings.Verbose = e.Verbose + replay.Settings.FileToReplyPath = "integration_request.gor" + replay.Settings.ForwardAddress = "127.0.0.1:" + strconv.Itoa(forwardPort) + replay.Settings.Port = port + + if e.ReplayLimit != 0 { + replay.Settings.ForwardAddress += "|" + strconv.Itoa(e.ReplayLimit) + } + + replay.Run() +} +func TestSavingRequestToFileAndReplyThem(t *testing.T) { + var request *http.Request + received := make(chan int) + + listenHandler := func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "OK", http.StatusNotFound) + } + + replayHandler := func(w http.ResponseWriter, r *http.Request) { + isEqual(t, r.URL.Path, request.URL.Path) + isEqual(t, r.Cookies()[0].Value, request.Cookies()[0].Value) + + http.Error(w, "404 page not found", http.StatusNotFound) + + if t.Failed() { + fmt.Println("\nReplayed:", r, "\nOriginal:", request) + } + + received <- 1 + } + + env := &Env{ + Verbose: true, + ListenHandler: listenHandler, + ReplayHandler: replayHandler, + } + p := env.startWithFile() + + request = getRequest(p) + + _, err := http.DefaultClient.Do(request) + + if err != nil { + t.Error("Can't make request", err) + } + + select { + case <-received: + case <-time.After(time.Second): + t.Error("Timeout error") + } +} From 6bb24697b2d831ea2e779c1ba3d29140999eccf2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 19 Aug 2013 19:00:08 +0200 Subject: [PATCH 03/31] add flags for file to listener --- listener/settings.go | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/listener/settings.go b/listener/settings.go index 4d6cd1c..e603d36 100644 --- a/listener/settings.go +++ b/listener/settings.go @@ -18,7 +18,8 @@ type ListenerSettings struct { Port int Address string - ReplayAddress string + ReplayAddress string + FileToReplyPath string ReplayLimit int @@ -48,5 +49,7 @@ func init() { flag.StringVar(&Settings.ReplayAddress, "r", defaultReplayAddress, "Address of replay server.") + flag.StringVar(&Settings.FileToReplyPath, "file", nil, "File to store captured requests") + flag.BoolVar(&Settings.Verbose, "verbose", false, "Log requests") } From 4f7b201ef110d0540bab8d5a4fcd1a0c615b2b23 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Mon, 19 Aug 2013 19:24:30 +0000 Subject: [PATCH 04/31] use empty string as default value --- listener/settings.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/listener/settings.go b/listener/settings.go index e603d36..4cf099c 100644 --- a/listener/settings.go +++ b/listener/settings.go @@ -49,7 +49,7 @@ func init() { flag.StringVar(&Settings.ReplayAddress, "r", defaultReplayAddress, "Address of replay server.") - flag.StringVar(&Settings.FileToReplyPath, "file", nil, "File to store captured requests") + flag.StringVar(&Settings.FileToReplyPath, "file", "", "File to store captured requests") flag.BoolVar(&Settings.Verbose, "verbose", false, "Log requests") } From e5cb5d959d4630787c697f805ba554724b533bd9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Mon, 19 Aug 2013 20:07:47 +0000 Subject: [PATCH 05/31] add flags for file in replay --- replay/settings.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/replay/settings.go b/replay/settings.go index 8ac8dfb..a45ca7d 100644 --- a/replay/settings.go +++ b/replay/settings.go @@ -19,6 +19,7 @@ type ReplaySettings struct { Host string ForwardAddress string + FileToReplyPath string Verbose bool } @@ -75,5 +76,7 @@ func init() { flag.StringVar(&Settings.ForwardAddress, "f", defaultAddress, "http address to forward traffic.\n\tYou can limit requests per second by adding `|num` after address.\n\tIf you have multiple addresses with different limits. For example: http://staging.example.com|100,http://dev.example.com|10") + flag.StringVar(&Settings.FileToReplyPath, "file", "", "File to replay captured requests from") + flag.BoolVar(&Settings.Verbose, "verbose", false, "Log requests") } From 5ccebd008d1567e7c8547b1695f563ac9468639a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Mon, 19 Aug 2013 20:08:38 +0000 Subject: [PATCH 06/31] test for saving listened requests to file wip --- listener/listener_test.go | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/listener/listener_test.go b/listener/listener_test.go index 319ed55..27ec4b6 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -43,3 +43,35 @@ func TestSendMessage(t *testing.T) { t.Errorf("Original and reveived requests does not match") } } + +func TestSaveMessageToFile(t *testing.T) { + Settings.Verbose = true + Settings.FileToReplyPath = "requests.gor" + + received := make(chan int) + Run() + + requestBytes = []byte("GET http://localhost:50000/pub/WWW/ HTTP/1.1\nHost: www.w3.org\r\n\r\n") + // TODO: implement foo + requestReader = foo(requestBytes) + request, err = http.ReadRequest(reader) + + go func() { + _, err := http.DefaultClient.Do(request) + received <- 1 + } + + select { + case <-received: + case <-time.After(time.Second): + t.Error("Timeout error") + } + + file, _ := os.Open("request.gor") + buf = make([]byte, 100) + n, _ = file.Read(buf) + + if bytes.Compare(buf, requestBytes) != 0 { + t.Errorf("Original and reveived requests does not match") + } +} From 95841d597ee83a181cd1f6e2fb9bdcff6b298a16 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 20 Aug 2013 09:05:41 +0200 Subject: [PATCH 07/31] wrap request buffer around reader --- listener/listener_test.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index 27ec4b6..40454ff 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -53,13 +53,13 @@ func TestSaveMessageToFile(t *testing.T) { requestBytes = []byte("GET http://localhost:50000/pub/WWW/ HTTP/1.1\nHost: www.w3.org\r\n\r\n") // TODO: implement foo - requestReader = foo(requestBytes) - request, err = http.ReadRequest(reader) + //requestReader = foo(requestBytes) + request, err = http.ReadRequest(bytes.NewBuffer(requestBytes)) go func() { _, err := http.DefaultClient.Do(request) received <- 1 - } + }() select { case <-received: From 700523e0c36328b26c23c79ff68a7b4c4b1808ae Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 20 Aug 2013 18:38:38 +0200 Subject: [PATCH 08/31] proper test for reader writing request to file --- listener/listener_test.go | 42 +++++++++++++++++++++++++++------------ 1 file changed, 29 insertions(+), 13 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index 40454ff..154bf67 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -4,6 +4,9 @@ import ( "bytes" "fmt" "net" + "net/http" + "os" + "time" "testing" ) @@ -13,6 +16,7 @@ func getTCPMessage() (msg *TCPMessage) { return &TCPMessage{packets: []*TCPPacket{packet}} } + func mockReplayServer() (listener net.Listener) { listener, _ = net.Listen("tcp", "127.0.0.1:0") @@ -35,7 +39,7 @@ func TestSendMessage(t *testing.T) { conn, _ := listener.Accept() defer conn.Close() - buf := make([]byte, 1024) + buf := make([]byte, 1024) n, _ := conn.Read(buf) buf = buf[0:n] @@ -49,16 +53,23 @@ func TestSaveMessageToFile(t *testing.T) { Settings.FileToReplyPath = "requests.gor" received := make(chan int) - Run() + go Run() - requestBytes = []byte("GET http://localhost:50000/pub/WWW/ HTTP/1.1\nHost: www.w3.org\r\n\r\n") - // TODO: implement foo - //requestReader = foo(requestBytes) - request, err = http.ReadRequest(bytes.NewBuffer(requestBytes)) + requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") + + handler := func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "OK", http.StatusNotFound) + received <- 1 + } + + go func() { + http.ListenAndServe(":50000", http.HandlerFunc(handler)) + time.Sleep(time.Millisecond * 100) + }() go func() { - _, err := http.DefaultClient.Do(request) - received <- 1 + conn, _ := net.Dial("tcp", ":50000") + conn.Write(requestBytes) }() select { @@ -67,11 +78,16 @@ func TestSaveMessageToFile(t *testing.T) { t.Error("Timeout error") } - file, _ := os.Open("request.gor") - buf = make([]byte, 100) - n, _ = file.Read(buf) + file, err := os.Open("request.gor") - if bytes.Compare(buf, requestBytes) != 0 { - t.Errorf("Original and reveived requests does not match") + if err != nil { + t.Errorf("Problem with opening file: ", err) + } + + fileBuf := make([]byte, 100) + file.Read(fileBuf) + + if bytes.Compare(fileBuf, requestBytes) != 0 { + t.Errorf("Original and received requests does not match") } } From ae792afa5eb780b57b1fbb341bd294e21416a8f0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 20 Aug 2013 18:40:27 +0200 Subject: [PATCH 09/31] gofmt test --- listener/listener_test.go | 45 +++++++++++++++++++-------------------- 1 file changed, 22 insertions(+), 23 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index 154bf67..e3f05b6 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -5,9 +5,9 @@ import ( "fmt" "net" "net/http" - "os" - "time" + "os" "testing" + "time" ) func getTCPMessage() (msg *TCPMessage) { @@ -16,7 +16,6 @@ func getTCPMessage() (msg *TCPMessage) { return &TCPMessage{packets: []*TCPPacket{packet}} } - func mockReplayServer() (listener net.Listener) { listener, _ = net.Listen("tcp", "127.0.0.1:0") @@ -39,7 +38,7 @@ func TestSendMessage(t *testing.T) { conn, _ := listener.Accept() defer conn.Close() - buf := make([]byte, 1024) + buf := make([]byte, 1024) n, _ := conn.Read(buf) buf = buf[0:n] @@ -50,27 +49,27 @@ func TestSendMessage(t *testing.T) { func TestSaveMessageToFile(t *testing.T) { Settings.Verbose = true - Settings.FileToReplyPath = "requests.gor" + Settings.FileToReplyPath = "requests.gor" received := make(chan int) - go Run() + go Run() - requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") + requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") - handler := func(w http.ResponseWriter, r *http.Request) { + handler := func(w http.ResponseWriter, r *http.Request) { http.Error(w, "OK", http.StatusNotFound) received <- 1 - } - - go func() { - http.ListenAndServe(":50000", http.HandlerFunc(handler)) - time.Sleep(time.Millisecond * 100) - }() + } go func() { - conn, _ := net.Dial("tcp", ":50000") - conn.Write(requestBytes) - }() + http.ListenAndServe(":50000", http.HandlerFunc(handler)) + time.Sleep(time.Millisecond * 100) + }() + + go func() { + conn, _ := net.Dial("tcp", ":50000") + conn.Write(requestBytes) + }() select { case <-received: @@ -78,14 +77,14 @@ func TestSaveMessageToFile(t *testing.T) { t.Error("Timeout error") } - file, err := os.Open("request.gor") + file, err := os.Open("request.gor") - if err != nil { - t.Errorf("Problem with opening file: ", err) - } + if err != nil { + t.Errorf("Problem with opening file: ", err) + } - fileBuf := make([]byte, 100) - file.Read(fileBuf) + fileBuf := make([]byte, 100) + file.Read(fileBuf) if bytes.Compare(fileBuf, requestBytes) != 0 { t.Errorf("Original and received requests does not match") From 54a3d73a9868c67189ca611ef4240dbb7ec79af7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 21 Aug 2013 18:29:21 +0200 Subject: [PATCH 10/31] better name --- listener/raw_tcp_listener.go | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/listener/raw_tcp_listener.go b/listener/raw_tcp_listener.go index ce2529b..f002192 100644 --- a/listener/raw_tcp_listener.go +++ b/listener/raw_tcp_listener.go @@ -25,18 +25,18 @@ type RAWTCPListener struct { port int // Port to listen } -func RAWTCPListen(addr string, port int) (listener *RAWTCPListener) { - listener = &RAWTCPListener{} +func RAWTCPListen(addr string, port int) (rawListener *RAWTCPListener) { + rawListener = &RAWTCPListener{} - listener.c_packets = make(chan *TCPPacket) - listener.c_messages = make(chan *TCPMessage) - listener.c_del_message = make(chan *TCPMessage) + rawListener.c_packets = make(chan *TCPPacket) + rawListener.c_messages = make(chan *TCPMessage) + rawListener.c_del_message = make(chan *TCPMessage) - listener.addr = addr - listener.port = port + rawListener.addr = addr + rawListener.port = port - go listener.listen() - go listener.readRAWSocket() + go rawListener.listen() + go rawListener.readRAWSocket() return } From db369cfc891860dfa180ec6f69105a395f1fc4a4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:10:07 +0200 Subject: [PATCH 11/31] gofmt --- replay/request_stats.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/replay/request_stats.go b/replay/request_stats.go index d1bbf32..86b5f61 100644 --- a/replay/request_stats.go +++ b/replay/request_stats.go @@ -43,8 +43,8 @@ func (s *RequestStat) IncResp(resp *HttpResponse) { // Updated stats timestamp to current time and reset to zero all stats values // TODO: Further on reset it should write stats to file func (s *RequestStat) reset() { - if s.timestamp != 0 { - Debug("Host:", s.host.Url, "Requests:", s.Count, "Errors:", s.Errors, "Status codes:", s.Codes) + if s.timestamp != 0 { + Debug("Host:", s.host.Url, "Requests:", s.Count, "Errors:", s.Errors, "Status codes:", s.Codes) } s.timestamp = time.Now().Unix() From 1e08ec39537c6f5e7cfa7ef863e3609a80dccf9b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:11:31 +0200 Subject: [PATCH 12/31] gofmt --- replay/settings.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/replay/settings.go b/replay/settings.go index a45ca7d..1f7ef04 100644 --- a/replay/settings.go +++ b/replay/settings.go @@ -18,7 +18,7 @@ type ReplaySettings struct { Port int Host string - ForwardAddress string + ForwardAddress string FileToReplyPath string Verbose bool From 02abcfdf744bc2bc1c15429d3477e03cbd0abc9f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:38:08 +0200 Subject: [PATCH 13/31] message logger --- listener/message_logger.go | 49 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 49 insertions(+) create mode 100644 listener/message_logger.go diff --git a/listener/message_logger.go b/listener/message_logger.go new file mode 100644 index 0000000..cdabcf1 --- /dev/null +++ b/listener/message_logger.go @@ -0,0 +1,49 @@ +package listener + +import ( + "fmt" + "os" +) + +type MessageLogger struct { + messageChannel chan string + + file *os.File +} + +func NewLog(filename string) *MessageLogger { + + file, err := os.OpenFile(filename, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0660) + + if err != nil { + panic(fmt.Sprintf("Cannot open file %q. Error: %s", filename, err)) + } + + messageLogger := &MessageLogger{ + messageChannel: make(chan string), + file: file, + } + + go func() { + defer func() { + messageLogger.close() + }() + + for { + select { + case message := <-messageLogger.messageChannel: + messageLogger.log(message) + } + } + }() + + return messageLogger +} + +func (messageLogger *MessageLogger) log(message string) { + fmt.Fprintln(messageLogger.file, message) +} + +func (messageLogger *MessageLogger) close() { + messageLogger.file.Close() +} From cfb65c06acf605a7077dc85f8c320deaf9c60fa2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:39:01 +0200 Subject: [PATCH 14/31] some temporary debug --- listener/tcp_message.go | 1 + 1 file changed, 1 insertion(+) diff --git a/listener/tcp_message.go b/listener/tcp_message.go index b34ec8c..90bd7f9 100644 --- a/listener/tcp_message.go +++ b/listener/tcp_message.go @@ -46,6 +46,7 @@ func (t *TCPMessage) listen() { for { select { case <-t.c_closing: + Debug("BLA BLA BLA") close(t.c_packets) return // Stop loop if message completed/expired case packet := <-t.c_packets: From f1eedc78a1e7a419ffa23abbb7c3c51b5781d7b0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:40:32 +0200 Subject: [PATCH 15/31] debugging and improving --- listener/listener_test.go | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index e3f05b6..41e233e 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -50,9 +50,10 @@ func TestSendMessage(t *testing.T) { func TestSaveMessageToFile(t *testing.T) { Settings.Verbose = true Settings.FileToReplyPath = "requests.gor" + Settings.Address = "127.0.0.1" + Settings.Port = 50000 received := make(chan int) - go Run() requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") @@ -63,9 +64,12 @@ func TestSaveMessageToFile(t *testing.T) { go func() { http.ListenAndServe(":50000", http.HandlerFunc(handler)) - time.Sleep(time.Millisecond * 100) }() + time.Sleep(time.Millisecond * 100) + time.Sleep(time.Millisecond * 100) + time.Sleep(time.Millisecond * 100) + go Run() go func() { conn, _ := net.Dial("tcp", ":50000") conn.Write(requestBytes) @@ -73,11 +77,12 @@ func TestSaveMessageToFile(t *testing.T) { select { case <-received: + time.Sleep(time.Millisecond * 100) case <-time.After(time.Second): t.Error("Timeout error") } - file, err := os.Open("request.gor") + file, err := os.Open("requests.gor") if err != nil { t.Errorf("Problem with opening file: ", err) From 335c86c6e9368d939ca0e18c7b77e5a6347ae79c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 28 Aug 2013 18:43:36 +0200 Subject: [PATCH 16/31] log messages --- listener/listener.go | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/listener/listener.go b/listener/listener.go index deeb671..f08d1c7 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -51,9 +51,17 @@ func Run() { currentTime := time.Now().UnixNano() currentRPS := 0 + var messageLogger *MessageLogger + + if Settings.FileToReplyPath != "" { + messageLogger = NewLog(Settings.FileToReplyPath) + } + for { // Receiving TCPMessage object + fmt.Println("FILE: ", Settings.FileToReplyPath) m := listener.Receive() + fmt.Println("bla bla bla: ", messageLogger) if Settings.ReplayLimit != 0 { if (time.Now().UnixNano() - currentTime) > time.Second.Nanoseconds() { @@ -68,6 +76,22 @@ func Run() { currentRPS++ } + fmt.Println("FOO BAR: ", messageLogger) + if messageLogger != nil { + go func() { + messageBuffer := new(bytes.Buffer) + messageWriter := bufio.NewWriter(messageBuffer) + + fmt.Fprintf(messageWriter, "\n") + fmt.Fprintf(messageWriter, "------------------------------------------------\n") + fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) + fmt.Fprintf(messageWriter, "\n") + + messageWriter.Flush() + messageLogger.messageChannel <- messageBuffer.String() + }() + } + fmt.Println("bla bla bla: ", messageLogger) go sendMessage(m) } } From 1a35a2c6511790cea55fc15f79e7ef8570bb7ffb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 23 Sep 2013 19:34:53 +0200 Subject: [PATCH 17/31] faulty listener test --- listener/listener_test.go | 76 ++++++++++++++++++++++++++------------- 1 file changed, 51 insertions(+), 25 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index 41e233e..2064dcd 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -7,7 +7,8 @@ import ( "net/http" "os" "testing" - "time" + "time" + "io" ) func getTCPMessage() (msg *TCPMessage) { @@ -21,8 +22,6 @@ func mockReplayServer() (listener net.Listener) { Settings.ReplayAddress = listener.Addr().String() - fmt.Println(listener.Addr().String()) - return } @@ -49,49 +48,76 @@ func TestSendMessage(t *testing.T) { func TestSaveMessageToFile(t *testing.T) { Settings.Verbose = true - Settings.FileToReplyPath = "requests.gor" + Settings.FileToReplyPath = "listener_test.gor" Settings.Address = "127.0.0.1" Settings.Port = 50000 - received := make(chan int) + receivedChan := make(chan int) - requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") + // receivedChan <- 1 + // requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") handler := func(w http.ResponseWriter, r *http.Request) { - http.Error(w, "OK", http.StatusNotFound) - received <- 1 + fmt.Println("handler called") + // fmt.Fprintf(w, "Hello, aaa") + // this is faulty + io.WriteString(w, "hello, world!\n") } + go Run() + go func() { http.ListenAndServe(":50000", http.HandlerFunc(handler)) }() - time.Sleep(time.Millisecond * 100) - time.Sleep(time.Millisecond * 100) - time.Sleep(time.Millisecond * 100) - go Run() - go func() { - conn, _ := net.Dial("tcp", ":50000") - conn.Write(requestBytes) - }() select { - case <-received: - time.Sleep(time.Millisecond * 100) + case msg, ok := <-receivedChan: + fmt.Println("received something") case <-time.After(time.Second): - t.Error("Timeout error") + fmt.Println("in timeout section") + // t.Error("Server not started and I dont know what is going on :(") } - file, err := os.Open("requests.gor") + request := getRequest() + resp, err := http.DefaultClient.Do(request) + if err != nil { + t.Errorf("Problem with default client", err) + } + fmt.Println("RESPONSE", resp) + + file, err := os.Open("listener_test.gor") if err != nil { t.Errorf("Problem with opening file: ", err) } - fileBuf := make([]byte, 100) - file.Read(fileBuf) + fileBuf := make([]byte, 1024) + n, err := file.Read(fileBuf) + fileBuf = fileBuf[:n] - if bytes.Compare(fileBuf, requestBytes) != 0 { - t.Errorf("Original and received requests does not match") - } + //requestBuffer := bytes.NewBuffer(fileBuf) + //requestReader := bufio.NewReader(requestBuffer) + //readRequest, _ := http.ReadRequest(requestReader) + fmt.Println("Read file: \n", string(fileBuf)) + fmt.Println("Read file: \n", fileBuf) + + //if bytes.Compare(fileBuf, make([]byte, 1024)) != 0 { + // t.Errorf("Original and received requests does not match") + //} + // if *request != *readRequest { + // t.Errorf("Original and received requests does not match") + //} + t.Errorf("Original and received requests does not match") +} + +func getRequest() (req *http.Request) { + req, _ = http.NewRequest("GET", "http://localhost:50000", nil) + ck := new(http.Cookie) + ck.Name = "test" + ck.Value = "value2" + + req.AddCookie(ck) + + return } From 5b6b8a53c1d3d8a046b4e6472557ea1617affcc2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 23 Sep 2013 19:36:23 +0200 Subject: [PATCH 18/31] better file formating --- listener/listener.go | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/listener/listener.go b/listener/listener.go index f08d1c7..6931a32 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -59,9 +59,7 @@ func Run() { for { // Receiving TCPMessage object - fmt.Println("FILE: ", Settings.FileToReplyPath) m := listener.Receive() - fmt.Println("bla bla bla: ", messageLogger) if Settings.ReplayLimit != 0 { if (time.Now().UnixNano() - currentTime) > time.Second.Nanoseconds() { @@ -76,22 +74,20 @@ func Run() { currentRPS++ } - fmt.Println("FOO BAR: ", messageLogger) if messageLogger != nil { + fmt.Println("FILE: ", Settings.FileToReplyPath) go func() { messageBuffer := new(bytes.Buffer) messageWriter := bufio.NewWriter(messageBuffer) - fmt.Fprintf(messageWriter, "\n") - fmt.Fprintf(messageWriter, "------------------------------------------------\n") + // fmt.Fprintf(messageWriter, "------------------------------------------------\n") fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) - fmt.Fprintf(messageWriter, "\n") messageWriter.Flush() messageLogger.messageChannel <- messageBuffer.String() }() } - fmt.Println("bla bla bla: ", messageLogger) + go sendMessage(m) } } From b562a6b62da201f047a05f6347735ceee888c234 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 23 Sep 2013 19:43:04 +0200 Subject: [PATCH 19/31] test is not working --- listener/listener_test.go | 82 +-------------------------------------- 1 file changed, 2 insertions(+), 80 deletions(-) diff --git a/listener/listener_test.go b/listener/listener_test.go index 2064dcd..319ed55 100644 --- a/listener/listener_test.go +++ b/listener/listener_test.go @@ -4,11 +4,7 @@ import ( "bytes" "fmt" "net" - "net/http" - "os" "testing" - "time" - "io" ) func getTCPMessage() (msg *TCPMessage) { @@ -22,6 +18,8 @@ func mockReplayServer() (listener net.Listener) { Settings.ReplayAddress = listener.Addr().String() + fmt.Println(listener.Addr().String()) + return } @@ -45,79 +43,3 @@ func TestSendMessage(t *testing.T) { t.Errorf("Original and reveived requests does not match") } } - -func TestSaveMessageToFile(t *testing.T) { - Settings.Verbose = true - Settings.FileToReplyPath = "listener_test.gor" - Settings.Address = "127.0.0.1" - Settings.Port = 50000 - - receivedChan := make(chan int) - - // receivedChan <- 1 - // requestBytes := []byte("GET / HTTP/1.1\nHost: localhost:50000\r\n\r\n") - - handler := func(w http.ResponseWriter, r *http.Request) { - fmt.Println("handler called") - // fmt.Fprintf(w, "Hello, aaa") - // this is faulty - io.WriteString(w, "hello, world!\n") - } - - go Run() - - go func() { - http.ListenAndServe(":50000", http.HandlerFunc(handler)) - }() - - - select { - case msg, ok := <-receivedChan: - fmt.Println("received something") - case <-time.After(time.Second): - fmt.Println("in timeout section") - // t.Error("Server not started and I dont know what is going on :(") - } - - request := getRequest() - resp, err := http.DefaultClient.Do(request) - if err != nil { - t.Errorf("Problem with default client", err) - } - fmt.Println("RESPONSE", resp) - - file, err := os.Open("listener_test.gor") - - if err != nil { - t.Errorf("Problem with opening file: ", err) - } - - fileBuf := make([]byte, 1024) - n, err := file.Read(fileBuf) - fileBuf = fileBuf[:n] - - //requestBuffer := bytes.NewBuffer(fileBuf) - //requestReader := bufio.NewReader(requestBuffer) - //readRequest, _ := http.ReadRequest(requestReader) - fmt.Println("Read file: \n", string(fileBuf)) - fmt.Println("Read file: \n", fileBuf) - - //if bytes.Compare(fileBuf, make([]byte, 1024)) != 0 { - // t.Errorf("Original and received requests does not match") - //} - // if *request != *readRequest { - // t.Errorf("Original and received requests does not match") - //} - t.Errorf("Original and received requests does not match") -} - -func getRequest() (req *http.Request) { - req, _ = http.NewRequest("GET", "http://localhost:50000", nil) - ck := new(http.Cookie) - ck.Name = "test" - ck.Value = "value2" - - req.AddCookie(ck) - - return -} From e7b2866ca9d54a72abc77b7f139a3911c95b5d7c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Mon, 23 Sep 2013 19:44:57 +0200 Subject: [PATCH 20/31] remove nasty debug --- listener/tcp_message.go | 1 - 1 file changed, 1 deletion(-) diff --git a/listener/tcp_message.go b/listener/tcp_message.go index 90bd7f9..b34ec8c 100644 --- a/listener/tcp_message.go +++ b/listener/tcp_message.go @@ -46,7 +46,6 @@ func (t *TCPMessage) listen() { for { select { case <-t.c_closing: - Debug("BLA BLA BLA") close(t.c_packets) return // Stop loop if message completed/expired case packet := <-t.c_packets: From 207eea1ba752f6175d3a5366bee180ec966d93f1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Sun, 29 Sep 2013 14:36:43 +0200 Subject: [PATCH 21/31] file replay partialy working --- integration_test.go | 51 ++++++++++++++++++++-------- listener/listener.go | 5 +-- replay/replay.go | 54 +++++++++++++++++++++++++++--- replay/replay_file_parser.go | 64 ++++++++++++++++++++++++++++++++++++ replay/request_factory.go | 3 +- 5 files changed, 154 insertions(+), 23 deletions(-) create mode 100644 replay/replay_file_parser.go diff --git a/integration_test.go b/integration_test.go index 795bbec..5fb383c 100644 --- a/integration_test.go +++ b/integration_test.go @@ -31,6 +31,7 @@ type Env struct { ReplayLimit int ListenerLimit int + ForwardPort int } func (e *Env) start() (p int) { @@ -199,13 +200,15 @@ func TestListenerRateLimit(t *testing.T) { } } -func (e *Env) startWithFile() (p int) { +func (e *Env) startFileListener() (p int) { p = 50000 + envs*10 + e.ForwardPort = p + 2 go e.startHTTP(p, http.HandlerFunc(e.ListenHandler)) go e.startHTTP(p+2, http.HandlerFunc(e.ReplayHandler)) go e.startFileUsingListener(p, p+1) - go e.startFileUsingReplay(p+1, p+2) + // we will replay after listener finishes capturing + // go e.startFileUsingReplay(p+1, p+2) // Time to start http and gor instances time.Sleep(time.Millisecond * 100) @@ -228,11 +231,10 @@ func (e *Env) startFileUsingListener(port int, replayPort int) { listener.Run() } -func (e *Env) startFileUsingReplay(port int, forwardPort int) { +func (e *Env) startFileUsingReplay() { replay.Settings.Verbose = e.Verbose replay.Settings.FileToReplyPath = "integration_request.gor" - replay.Settings.ForwardAddress = "127.0.0.1:" + strconv.Itoa(forwardPort) - replay.Settings.Port = port + replay.Settings.ForwardAddress = "127.0.0.1:" + strconv.Itoa(e.ForwardPort) if e.ReplayLimit != 0 { replay.Settings.ForwardAddress += "|" + strconv.Itoa(e.ReplayLimit) @@ -240,25 +242,33 @@ func (e *Env) startFileUsingReplay(port int, forwardPort int) { replay.Run() } + func TestSavingRequestToFileAndReplyThem(t *testing.T) { var request *http.Request - received := make(chan int) + processed := make(chan int) listenHandler := func(w http.ResponseWriter, r *http.Request) { http.Error(w, "OK", http.StatusNotFound) } + requestsCount := 0 + var replayedRequests []*http.Request replayHandler := func(w http.ResponseWriter, r *http.Request) { + requestsCount++ + isEqual(t, r.URL.Path, request.URL.Path) isEqual(t, r.Cookies()[0].Value, request.Cookies()[0].Value) http.Error(w, "404 page not found", http.StatusNotFound) + replayedRequests = append(replayedRequests, r) if t.Failed() { fmt.Println("\nReplayed:", r, "\nOriginal:", request) } - received <- 1 + if requestsCount > 1 { + processed <- 1 + } } env := &Env{ @@ -266,19 +276,32 @@ func TestSavingRequestToFileAndReplyThem(t *testing.T) { ListenHandler: listenHandler, ReplayHandler: replayHandler, } - p := env.startWithFile() + + p := env.startFileListener() request = getRequest(p) - _, err := http.DefaultClient.Do(request) + for i := 0; i < 2; i++ { + go func() { + _, err := http.DefaultClient.Do(request) + + if err != nil { + t.Error("Can't make request", err) + } + }() + } + + // TODO: wait until gor will process response, should be kind of flag/semaphore + time.Sleep(time.Millisecond * 700) + go env.startFileUsingReplay() - if err != nil { - t.Error("Can't make request", err) - } select { - case <-received: - case <-time.After(time.Second): + case <-processed: + case <-time.After(2 * time.Second): + for _, value := range replayedRequests { + fmt.Println(value) + } t.Error("Timeout error") } } diff --git a/listener/listener.go b/listener/listener.go index 6931a32..fc533c8 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -80,15 +80,16 @@ func Run() { messageBuffer := new(bytes.Buffer) messageWriter := bufio.NewWriter(messageBuffer) - // fmt.Fprintf(messageWriter, "------------------------------------------------\n") fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) + fmt.Fprintf(messageWriter, "\n--\n") messageWriter.Flush() messageLogger.messageChannel <- messageBuffer.String() }() + } else { + go sendMessage(m) } - go sendMessage(m) } } diff --git a/replay/replay.go b/replay/replay.go index efb487b..1a18a3e 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -35,6 +35,17 @@ import ( const bufSize = 4096 +type ReplayManager struct { + reqFactory *RequestFactory +} + +func NewReplayManager() (rm *ReplayManager) { + rm = &ReplayManager{} + rm.reqFactory = NewRequestFactory() + + return +} + // Enable debug logging only if "--verbose" flag passed func Debug(v ...interface{}) { if Settings.Verbose { @@ -54,6 +65,37 @@ func ParseRequest(data []byte) (request *http.Request, err error) { // Replay server listen to UDP traffic from Listeners // Each request processed by RequestFactory func Run() { + rm := NewReplayManager(); + + if Settings.FileToReplyPath != "" { + rm.RunReplayFromFile() + } else { + rm.RunReplayFromNetwork() + } + +} + +func (self *ReplayManager) RunReplayFromFile() { + log.Println("Starting file reply") + requests, err := parseReplyFile() + if err != nil { + log.Fatal("Can't parse request:", err) + } + + for _, request := range requests { + parsedReq, err := ParseRequest(request) + if err != nil { + log.Fatal("Can't parse request:", err) + } + + self.sendRequestToReplay(parsedReq) + } + + // wait forever + select {} +} + +func (self *ReplayManager) RunReplayFromNetwork() { listener, err := net.Listen("tcp", Settings.Address()) log.Println("Starting replay server at:", Settings.Address()) @@ -66,8 +108,6 @@ func Run() { log.Println("Forwarding requests to:", host.Url, "limit:", host.Limit) } - requestFactory := NewRequestFactory() - for { conn, err := listener.Accept() @@ -76,12 +116,12 @@ func Run() { continue } - go handleConnection(conn, requestFactory) + go self.handleConnection(conn) } } -func handleConnection(conn net.Conn, rf *RequestFactory) error { +func (self *ReplayManager) handleConnection(conn net.Conn) error { defer conn.Close() var read = true @@ -112,9 +152,13 @@ func handleConnection(conn net.Conn, rf *RequestFactory) error { } else { Debug("Adding request", request) - rf.Add(request) + self.sendRequestToReplay(request) } }() return nil } + +func (self *ReplayManager) sendRequestToReplay(req *http.Request) { + self.reqFactory.Add(req) +} diff --git a/replay/replay_file_parser.go b/replay/replay_file_parser.go new file mode 100644 index 0000000..2068e37 --- /dev/null +++ b/replay/replay_file_parser.go @@ -0,0 +1,64 @@ +package replay + +import ( + "bufio" + "log" + "os" + "bytes" +) + +func parseReplyFile() (requests [][]byte, err error) { + requests, err = readLines(Settings.FileToReplyPath) + + if err != nil { + log.Fatalf("readLines: %s", err) + } + + return +} + +// readLines reads a whole file into memory +// and returns a slice of its lines. +func readLines(path string) (requests [][]byte, err error) { + file, err := os.Open(path) + if err != nil { + return nil, err + } + defer file.Close() + + scanner := bufio.NewScanner(file) + scanner.Split(scanLinesFunc) + for scanner.Scan() { + if len(scanner.Text()) > 5 { + requests = append(requests, scanner.Bytes()) + } + } + return requests, scanner.Err() +} + +// scanner spliting logic +func scanLinesFunc(data []byte, atEOF bool) (advance int, token []byte, err error) { + if atEOF && len(data) == 0 { + return 0, nil, nil + } + + delimiter := []byte{'\n','-','-','\n'} + if i := bytes.Index(data, delimiter); i >= 0 { + // We have a full newline-terminated line. + return i + len(delimiter), dropCR(data[0:i]), nil + } + // If we're at EOF, we have a final, non-terminated line. Return it. + if atEOF { + return len(data), dropCR(data), nil + } + // Request more data. + return 0, nil, nil +} + +func dropCR(data []byte) []byte { + if len(data) > 0 && data[len(data)-1] == '\r' { + return data[0 : len(data)-1] + } + return data +} + diff --git a/replay/request_factory.go b/replay/request_factory.go index 1d8e65f..77a5b3a 100644 --- a/replay/request_factory.go +++ b/replay/request_factory.go @@ -58,8 +58,6 @@ func (f *RequestFactory) sendRequest(host *ForwardHost, request *http.Request) { request.RequestURI = "" request.URL, _ = url.ParseRequestURI(URL) - Debug("Sending request:", host.Url, request) - resp, err := client.Do(request) if err == nil { @@ -94,6 +92,7 @@ func (f *RequestFactory) handleRequests() { resp.host.Stat.IncResp(resp) } } + } // Add request to channel for further processing From 47b498d4d961345e347fd7da213bdef00132bd97 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 2 Oct 2013 19:27:52 +0200 Subject: [PATCH 22/31] fix request reading --- integration_test.go | 47 +++++++------ listener/listener.go | 2 +- replay/replay.go | 129 ++++++++++++++++++----------------- replay/replay_file_parser.go | 24 +++---- 4 files changed, 104 insertions(+), 98 deletions(-) diff --git a/integration_test.go b/integration_test.go index 5fb383c..77dba65 100644 --- a/integration_test.go +++ b/integration_test.go @@ -31,7 +31,7 @@ type Env struct { ReplayLimit int ListenerLimit int - ForwardPort int + ForwardPort int } func (e *Env) start() (p int) { @@ -203,11 +203,11 @@ func TestListenerRateLimit(t *testing.T) { func (e *Env) startFileListener() (p int) { p = 50000 + envs*10 - e.ForwardPort = p + 2 + e.ForwardPort = p + 2 go e.startHTTP(p, http.HandlerFunc(e.ListenHandler)) go e.startHTTP(p+2, http.HandlerFunc(e.ReplayHandler)) go e.startFileUsingListener(p, p+1) - // we will replay after listener finishes capturing + // we will replay after listener finishes capturing // go e.startFileUsingReplay(p+1, p+2) // Time to start http and gor instances @@ -251,24 +251,24 @@ func TestSavingRequestToFileAndReplyThem(t *testing.T) { http.Error(w, "OK", http.StatusNotFound) } - requestsCount := 0 - var replayedRequests []*http.Request + requestsCount := 0 + var replayedRequests []*http.Request replayHandler := func(w http.ResponseWriter, r *http.Request) { - requestsCount++ + requestsCount++ isEqual(t, r.URL.Path, request.URL.Path) isEqual(t, r.Cookies()[0].Value, request.Cookies()[0].Value) http.Error(w, "404 page not found", http.StatusNotFound) - replayedRequests = append(replayedRequests, r) + replayedRequests = append(replayedRequests, r) if t.Failed() { fmt.Println("\nReplayed:", r, "\nOriginal:", request) } - if requestsCount > 1 { - processed <- 1 - } + if requestsCount > 1 { + processed <- 1 + } } env := &Env{ @@ -281,27 +281,26 @@ func TestSavingRequestToFileAndReplyThem(t *testing.T) { request = getRequest(p) - for i := 0; i < 2; i++ { - go func() { - _, err := http.DefaultClient.Do(request) + for i := 0; i < 2; i++ { + go func() { + _, err := http.DefaultClient.Do(request) - if err != nil { - t.Error("Can't make request", err) - } - }() - } + if err != nil { + t.Error("Can't make request", err) + } + }() + } - // TODO: wait until gor will process response, should be kind of flag/semaphore + // TODO: wait until gor will process response, should be kind of flag/semaphore time.Sleep(time.Millisecond * 700) - go env.startFileUsingReplay() - + go env.startFileUsingReplay() select { case <-processed: case <-time.After(2 * time.Second): - for _, value := range replayedRequests { - fmt.Println(value) - } + for _, value := range replayedRequests { + fmt.Println(value) + } t.Error("Timeout error") } } diff --git a/listener/listener.go b/listener/listener.go index fc533c8..a6ceff0 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -81,7 +81,7 @@ func Run() { messageWriter := bufio.NewWriter(messageBuffer) fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) - fmt.Fprintf(messageWriter, "\n--\n") + // fmt.Fprintf(messageWriter, "\n--\n") messageWriter.Flush() messageLogger.messageChannel <- messageBuffer.String() diff --git a/replay/replay.go b/replay/replay.go index 1a18a3e..acca7c5 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -25,12 +25,12 @@ package replay import ( - "bufio" - "bytes" - "io" - "log" - "net" - "net/http" + "bufio" + "bytes" + "io" + "log" + "net" + "net/http" ) const bufSize = 4096 @@ -41,24 +41,29 @@ type ReplayManager struct { func NewReplayManager() (rm *ReplayManager) { rm = &ReplayManager{} - rm.reqFactory = NewRequestFactory() + rm.reqFactory = NewRequestFactory() return } // Enable debug logging only if "--verbose" flag passed func Debug(v ...interface{}) { - if Settings.Verbose { - log.Println(v...) - } + if Settings.Verbose { + log.Println(v...) + } } func ParseRequest(data []byte) (request *http.Request, err error) { - buf := bytes.NewBuffer(data) - reader := bufio.NewReader(buf) + buf := bytes.NewBuffer(data) + reader := bufio.NewReader(buf) - request, err = http.ReadRequest(reader) - return + request, err = http.ReadRequest(reader) + + if err != nil { + log.Fatal("Can not parse request", string(data), err) + } + + return } // Because its sub-program, Run acts as `main` @@ -76,16 +81,18 @@ func Run() { } func (self *ReplayManager) RunReplayFromFile() { - log.Println("Starting file reply") + log.Println("Starting file reply") requests, err := parseReplyFile() + if err != nil { - log.Fatal("Can't parse request:", err) + log.Fatal("Can't parse request: ", err) } for _, request := range requests { parsedReq, err := ParseRequest(request) + if err != nil { - log.Fatal("Can't parse request:", err) + log.Fatal("Can't parse request...:", err) } self.sendRequestToReplay(parsedReq) @@ -96,67 +103,67 @@ func (self *ReplayManager) RunReplayFromFile() { } func (self *ReplayManager) RunReplayFromNetwork() { - listener, err := net.Listen("tcp", Settings.Address()) + listener, err := net.Listen("tcp", Settings.Address()) - log.Println("Starting replay server at:", Settings.Address()) + log.Println("Starting replay server at:", Settings.Address()) - if err != nil { - log.Fatal("Can't start:", err) - } + if err != nil { + log.Fatal("Can't start:", err) + } - for _, host := range Settings.ForwardedHosts() { - log.Println("Forwarding requests to:", host.Url, "limit:", host.Limit) - } + for _, host := range Settings.ForwardedHosts() { + log.Println("Forwarding requests to:", host.Url, "limit:", host.Limit) + } - for { - conn, err := listener.Accept() + for { + conn, err := listener.Accept() - if err != nil { - log.Println("Error while Accept()", err) - continue - } + if err != nil { + log.Println("Error while Accept()", err) + continue + } - go self.handleConnection(conn) - } + go self.handleConnection(conn) + } } func (self *ReplayManager) handleConnection(conn net.Conn) error { - defer conn.Close() + defer conn.Close() - var read = true - var response []byte - var buf []byte + var read = true + var response []byte + var buf []byte - buf = make([]byte, bufSize) + buf = make([]byte, bufSize) - for read { - n, err := conn.Read(buf) + for read { + n, err := conn.Read(buf) - switch err { - case io.EOF: - read = false - case nil: - response = append(response, buf[0:n]...) - if n < bufSize { - read = false - } - default: - read = false - } - } + switch err { + case io.EOF: + read = false + case nil: + response = append(response, buf[0:n]...) + if n < bufSize { + read = false + } + default: + read = false + } + } - go func() { - if request, err := ParseRequest(response); err != nil { - Debug("Error while parsing request", err, response) - } else { - Debug("Adding request", request) + go func() { + if request, err := ParseRequest(response); err != nil { + Debug("Error while parsing request", err, response) + } else { + Debug("Adding request", request) - self.sendRequestToReplay(request) - } - }() + self.sendRequestToReplay(request) + } + }() - return nil + return nil } func (self *ReplayManager) sendRequestToReplay(req *http.Request) { diff --git a/replay/replay_file_parser.go b/replay/replay_file_parser.go index 2068e37..6454ce8 100644 --- a/replay/replay_file_parser.go +++ b/replay/replay_file_parser.go @@ -21,6 +21,7 @@ func parseReplyFile() (requests [][]byte, err error) { // and returns a slice of its lines. func readLines(path string) (requests [][]byte, err error) { file, err := os.Open(path) + if err != nil { return nil, err } @@ -28,11 +29,13 @@ func readLines(path string) (requests [][]byte, err error) { scanner := bufio.NewScanner(file) scanner.Split(scanLinesFunc) + for scanner.Scan() { if len(scanner.Text()) > 5 { requests = append(requests, scanner.Bytes()) } } + return requests, scanner.Err() } @@ -42,23 +45,20 @@ func scanLinesFunc(data []byte, atEOF bool) (advance int, token []byte, err erro return 0, nil, nil } - delimiter := []byte{'\n','-','-','\n'} + delimiter := []byte{'\r', '\n', '\r', '\n', '\n'} + if i := bytes.Index(data, delimiter); i >= 0 { - // We have a full newline-terminated line. - return i + len(delimiter), dropCR(data[0:i]), nil + // We have a http request end: \r\n\r\n + log.Printf("to read: %v", i + len(delimiter)) + + return (i + len(delimiter)), data[0:(i + len(delimiter))], nil } + // If we're at EOF, we have a final, non-terminated line. Return it. if atEOF { - return len(data), dropCR(data), nil + return len(data), data, nil } + // Request more data. return 0, nil, nil } - -func dropCR(data []byte) []byte { - if len(data) > 0 && data[len(data)-1] == '\r' { - return data[0 : len(data)-1] - } - return data -} - From 20ea55e6e7d55d5550b08cc0f401daef96e11437 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 2 Oct 2013 19:34:54 +0200 Subject: [PATCH 23/31] remove debug statement --- replay/replay_file_parser.go | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/replay/replay_file_parser.go b/replay/replay_file_parser.go index 6454ce8..6910ff5 100644 --- a/replay/replay_file_parser.go +++ b/replay/replay_file_parser.go @@ -47,10 +47,8 @@ func scanLinesFunc(data []byte, atEOF bool) (advance int, token []byte, err erro delimiter := []byte{'\r', '\n', '\r', '\n', '\n'} + // We have a http request end: \r\n\r\n if i := bytes.Index(data, delimiter); i >= 0 { - // We have a http request end: \r\n\r\n - log.Printf("to read: %v", i + len(delimiter)) - return (i + len(delimiter)), data[0:(i + len(delimiter))], nil } From 30547f2a4174eef1862238bc3e4ffe103d6458fc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Wed, 2 Oct 2013 20:07:19 +0200 Subject: [PATCH 24/31] Embed log.Logger in MessageLogger --- listener/listener.go | 4 ++-- listener/message_logger.go | 28 ++++++---------------------- replay/replay.go | 1 - 3 files changed, 8 insertions(+), 25 deletions(-) diff --git a/listener/listener.go b/listener/listener.go index 8d06f26..bc547ca 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -81,11 +81,11 @@ func Run() { messageBuffer := new(bytes.Buffer) messageWriter := bufio.NewWriter(messageBuffer) + // TODO: add timestamp to message fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) - // fmt.Fprintf(messageWriter, "\n--\n") messageWriter.Flush() - messageLogger.messageChannel <- messageBuffer.String() + messageLogger.Println(messageBuffer.String()) }() } else { go sendMessage(m) diff --git a/listener/message_logger.go b/listener/message_logger.go index cdabcf1..c4444e6 100644 --- a/listener/message_logger.go +++ b/listener/message_logger.go @@ -1,12 +1,12 @@ package listener import ( - "fmt" + "log" "os" ) type MessageLogger struct { - messageChannel chan string + *log.Logger file *os.File } @@ -16,34 +16,18 @@ func NewLog(filename string) *MessageLogger { file, err := os.OpenFile(filename, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0660) if err != nil { - panic(fmt.Sprintf("Cannot open file %q. Error: %s", filename, err)) + log.Fatal("Cannot open file %q. Error: %s", filename, err) } + logger := log.New(file, "", 0) + messageLogger := &MessageLogger{ - messageChannel: make(chan string), - file: file, + Logger: logger, } - go func() { - defer func() { - messageLogger.close() - }() - - for { - select { - case message := <-messageLogger.messageChannel: - messageLogger.log(message) - } - } - }() - return messageLogger } -func (messageLogger *MessageLogger) log(message string) { - fmt.Fprintln(messageLogger.file, message) -} - func (messageLogger *MessageLogger) close() { messageLogger.file.Close() } diff --git a/replay/replay.go b/replay/replay.go index 2eace6d..bacde72 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -125,7 +125,6 @@ func (self *ReplayManager) RunReplayFromNetwork() { go self.handleConnection(conn) } - } func (self *ReplayManager) handleConnection(conn net.Conn) error { From 8b66a99665411b4ba8536f30a17317a933782fac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Thu, 3 Oct 2013 21:44:26 +0000 Subject: [PATCH 25/31] Save request timestamp WIP --- listener/listener.go | 2 ++ replay/replay.go | 14 +++++++++++++- replay/replay_file_parser.go | 24 +++++++++++++++++++++--- 3 files changed, 36 insertions(+), 4 deletions(-) diff --git a/listener/listener.go b/listener/listener.go index bc547ca..9b72c94 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -82,6 +82,8 @@ func Run() { messageWriter := bufio.NewWriter(messageBuffer) // TODO: add timestamp to message + fmt.Fprintf(messageWriter, "%v\n", time.Now().UnixNano()) + fmt.Printf("timestamp: %v", time.Now().UnixNano()) fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) messageWriter.Flush() diff --git a/replay/replay.go b/replay/replay.go index bacde72..649cf51 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -31,8 +31,10 @@ import ( "log" "net" "net/http" + "time" ) + const bufSize = 4096 type ReplayManager struct { @@ -88,17 +90,27 @@ func (self *ReplayManager) RunReplayFromFile() { log.Fatal("Can't parse request: ", err) } + var lastTimestamp int64 + + if len(requests) > 0 { + lastTimestamp = requests[0].Timestamp + } for _, request := range requests { - parsedReq, err := ParseRequest(request) + + parsedReq, err := ParseRequest(request.Request) if err != nil { log.Fatal("Can't parse request...:", err) } + time.Sleep(time.Duration(request.Timestamp - lastTimestamp)) + self.sendRequestToReplay(parsedReq) + lastTimestamp = request.Timestamp } // wait forever + // TODO: quit when all request finishes select {} } diff --git a/replay/replay_file_parser.go b/replay/replay_file_parser.go index 6910ff5..81c82df 100644 --- a/replay/replay_file_parser.go +++ b/replay/replay_file_parser.go @@ -5,9 +5,21 @@ import ( "log" "os" "bytes" + "strconv" + + "fmt" ) -func parseReplyFile() (requests [][]byte, err error) { +type ParsedRequest struct { + Request []byte + Timestamp int64 +} + +func (self ParsedRequest) String() string { + return fmt.Sprintf("Request: %v, timestamp: %v", string(self.Request), self.Timestamp) +} + +func parseReplyFile() (requests []ParsedRequest, err error) { requests, err = readLines(Settings.FileToReplyPath) if err != nil { @@ -19,7 +31,7 @@ func parseReplyFile() (requests [][]byte, err error) { // readLines reads a whole file into memory // and returns a slice of its lines. -func readLines(path string) (requests [][]byte, err error) { +func readLines(path string) (requests []ParsedRequest, err error) { file, err := os.Open(path) if err != nil { @@ -32,7 +44,13 @@ func readLines(path string) (requests [][]byte, err error) { for scanner.Scan() { if len(scanner.Text()) > 5 { - requests = append(requests, scanner.Bytes()) + Debug(scanner.Text()) + i := bytes.IndexByte(scanner.Bytes(), '\n') + timestamp, _ := strconv.Atoi(string(scanner.Bytes()[i + 1:])) + pr := ParsedRequest{scanner.Bytes()[:i], int64(timestamp)} + + Debug(pr) + requests = append(requests, pr) } } From 5782d486f9ea69dfa096aa6fcbe698eb76c08188 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz=20S=C5=82awi=C5=84ski?= Date: Fri, 4 Oct 2013 04:56:48 +0000 Subject: [PATCH 26/31] Parse request with timestamp properly --- listener/listener.go | 1 - replay/replay_file_parser.go | 6 ++---- 2 files changed, 2 insertions(+), 5 deletions(-) diff --git a/listener/listener.go b/listener/listener.go index 9b72c94..10a2e9c 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -83,7 +83,6 @@ func Run() { // TODO: add timestamp to message fmt.Fprintf(messageWriter, "%v\n", time.Now().UnixNano()) - fmt.Printf("timestamp: %v", time.Now().UnixNano()) fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) messageWriter.Flush() diff --git a/replay/replay_file_parser.go b/replay/replay_file_parser.go index 81c82df..a2fac60 100644 --- a/replay/replay_file_parser.go +++ b/replay/replay_file_parser.go @@ -44,12 +44,10 @@ func readLines(path string) (requests []ParsedRequest, err error) { for scanner.Scan() { if len(scanner.Text()) > 5 { - Debug(scanner.Text()) i := bytes.IndexByte(scanner.Bytes(), '\n') - timestamp, _ := strconv.Atoi(string(scanner.Bytes()[i + 1:])) - pr := ParsedRequest{scanner.Bytes()[:i], int64(timestamp)} + timestamp, _ := strconv.Atoi(string(scanner.Bytes()[:i])) + pr := ParsedRequest{scanner.Bytes()[i + 1:], int64(timestamp)} - Debug(pr) requests = append(requests, pr) } } From 830ae237c9e5d92e81f3a845871819b58d225274 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 8 Oct 2013 12:54:54 +0200 Subject: [PATCH 27/31] better messages on listener start --- listener/listener.go | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/listener/listener.go b/listener/listener.go index 10a2e9c..f96286c 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -44,7 +44,18 @@ func Run() { } fmt.Println("Listening for HTTP traffic on", Settings.Address+":"+strconv.Itoa(Settings.Port)) - fmt.Println("Forwarding requests to replay server:", Settings.ReplayAddress, "Limit:", Settings.ReplayLimit) + + var messageLogger *MessageLogger + + if Settings.FileToReplyPath != "" { + messageLogger = NewLog(Settings.FileToReplyPath) + } + + if messageLogger == nil { + fmt.Println("Forwarding requests to replay server:", Settings.ReplayAddress, "Limit:", Settings.ReplayLimit) + } else { + fmt.Println("Saving requests to file", Settings.FileToReplyPath) + } // Sniffing traffic from given address listener := RAWTCPListen(Settings.Address, Settings.Port) @@ -52,11 +63,6 @@ func Run() { currentTime := time.Now().UnixNano() currentRPS := 0 - var messageLogger *MessageLogger - - if Settings.FileToReplyPath != "" { - messageLogger = NewLog(Settings.FileToReplyPath) - } for { // Receiving TCPMessage object @@ -76,12 +82,10 @@ func Run() { } if messageLogger != nil { - fmt.Println("FILE: ", Settings.FileToReplyPath) go func() { messageBuffer := new(bytes.Buffer) messageWriter := bufio.NewWriter(messageBuffer) - // TODO: add timestamp to message fmt.Fprintf(messageWriter, "%v\n", time.Now().UnixNano()) fmt.Fprintf(messageWriter, "%s", string(m.Bytes())) From 9077656bdfb77c23de01cc905c55398d85a03a79 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 8 Oct 2013 12:59:11 +0200 Subject: [PATCH 28/31] Information about replaying in README --- README.md | 22 ++++++++++++++++------ 1 file changed, 16 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index 96fb0c4..45cbf37 100644 --- a/README.md +++ b/README.md @@ -2,17 +2,17 @@ ## About -Gor is a simple http traffic replication tool written in Go. +Gor is a simple http traffic replication tool written in Go. Its main goal is to replay traffic from production servers to staging and dev environments. -Now you can test your code on real user sessions in an automated and repeatable fashion. +Now you can test your code on real user sessions in an automated and repeatable fashion. **No more falling down in production!** Gor consists of 2 parts: listener and replay servers. The listener server catches http traffic from a given port in real-time -and sends it to the replay server. +and sends it to the replay server or saves to file. The replay server forwards traffic to a given address. @@ -23,9 +23,9 @@ The replay server forwards traffic to a given address. ```bash # Run on servers where you want to catch traffic. You can run it on each `web` machine. -sudo gor listen -p 80 -r replay.server.local:28020 +sudo gor listen -p 80 -r replay.server.local:28020 -# Replay server (replay.server.local). +# Replay server (replay.server.local). gor replay -f http://staging.server -p 28020 ``` @@ -54,6 +54,16 @@ You can forward traffic to multiple endpoints. Just separate the addresses by co ``` gor replay -f "http://staging.server|10,http://dev.server|5" ``` +### Saving requests to file +You can save request to save to file to replay multiple time, or in different network: +``` +gor listen -p 8080 -file requests.gor +``` + +And replaying: +``` +gor replay -f "http://staging.server" -file requests.gor +``` ## Additional help ``` @@ -80,7 +90,7 @@ https://github.com/buger/gor/releases ## Building from source 1. Setup standard Go environment http://golang.org/doc/code.html and ensure that $GOPATH environment variable properly set. -2. `go get github.com/buger/gor`. +2. `go get github.com/buger/gor`. 3. `cd $GOPATH/src/github.com/buger/gor` 4. `go build gor.go` to get binary, or `go run gor.go` to build and run (useful for development) From 8adeda9a9888a9a24275295bc51a17def862245f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 8 Oct 2013 19:41:24 +0200 Subject: [PATCH 29/31] finish replay, when all requests got replayed --- replay/replay.go | 23 ++++++++++++++++++++--- replay/request_stats.go | 3 +++ 2 files changed, 23 insertions(+), 3 deletions(-) diff --git a/replay/replay.go b/replay/replay.go index 649cf51..6fdd5e1 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -83,6 +83,8 @@ func Run() { } func (self *ReplayManager) RunReplayFromFile() { + TotalResponsesCount = 0 + log.Println("Starting file reply") requests, err := parseReplyFile() @@ -95,6 +97,20 @@ func (self *ReplayManager) RunReplayFromFile() { if len(requests) > 0 { lastTimestamp = requests[0].Timestamp } + + requestsToReplay := 0 + + hosts := Settings.ForwardedHosts() + log.Println("requests", len(requests)) + for _, host := range hosts { + log.Println(host.Limit) + if host.Limit > 0 { + requestsToReplay += host.Limit + } else { + requestsToReplay += len(requests) + } + } + for _, request := range requests { parsedReq, err := ParseRequest(request.Request) @@ -109,9 +125,10 @@ func (self *ReplayManager) RunReplayFromFile() { lastTimestamp = request.Timestamp } - // wait forever - // TODO: quit when all request finishes - select {} + for requestsToReplay > TotalResponsesCount { + time.Sleep(time.Second) + } + } func (self *ReplayManager) RunReplayFromNetwork() { diff --git a/replay/request_stats.go b/replay/request_stats.go index cf26282..149eb58 100644 --- a/replay/request_stats.go +++ b/replay/request_stats.go @@ -4,6 +4,8 @@ import ( "time" ) +var TotalResponsesCount int + // RequestStat stores in context of current timestamp type RequestStat struct { timestamp int64 @@ -31,6 +33,7 @@ func (s *RequestStat) IncReq() { // IncResp is called after response func (s *RequestStat) IncResp(resp *HttpResponse) { s.Touch() + TotalResponsesCount++ if resp.err != nil { s.Errors++ From 685c998e8ed8d1009adb0718477a08f821ee8e6f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 8 Oct 2013 19:47:30 +0200 Subject: [PATCH 30/31] remove MessageLogger --- listener/listener.go | 13 ++++++++++--- listener/message_logger.go | 33 --------------------------------- replay/replay.go | 2 -- 3 files changed, 10 insertions(+), 38 deletions(-) delete mode 100644 listener/message_logger.go diff --git a/listener/listener.go b/listener/listener.go index f96286c..5cf1e79 100644 --- a/listener/listener.go +++ b/listener/listener.go @@ -45,10 +45,18 @@ func Run() { fmt.Println("Listening for HTTP traffic on", Settings.Address+":"+strconv.Itoa(Settings.Port)) - var messageLogger *MessageLogger + var messageLogger *log.Logger if Settings.FileToReplyPath != "" { - messageLogger = NewLog(Settings.FileToReplyPath) + + file, err := os.OpenFile(Settings.FileToReplyPath, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0660) + defer file.Close() + + if err != nil { + log.Fatal("Cannot open file %q. Error: %s", Settings.FileToReplyPath, err) + } + + messageLogger = log.New(file, "", 0) } if messageLogger == nil { @@ -95,7 +103,6 @@ func Run() { } else { go sendMessage(m) } - } } diff --git a/listener/message_logger.go b/listener/message_logger.go deleted file mode 100644 index c4444e6..0000000 --- a/listener/message_logger.go +++ /dev/null @@ -1,33 +0,0 @@ -package listener - -import ( - "log" - "os" -) - -type MessageLogger struct { - *log.Logger - - file *os.File -} - -func NewLog(filename string) *MessageLogger { - - file, err := os.OpenFile(filename, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0660) - - if err != nil { - log.Fatal("Cannot open file %q. Error: %s", filename, err) - } - - logger := log.New(file, "", 0) - - messageLogger := &MessageLogger{ - Logger: logger, - } - - return messageLogger -} - -func (messageLogger *MessageLogger) close() { - messageLogger.file.Close() -} diff --git a/replay/replay.go b/replay/replay.go index 6fdd5e1..62099d7 100644 --- a/replay/replay.go +++ b/replay/replay.go @@ -101,9 +101,7 @@ func (self *ReplayManager) RunReplayFromFile() { requestsToReplay := 0 hosts := Settings.ForwardedHosts() - log.Println("requests", len(requests)) for _, host := range hosts { - log.Println(host.Limit) if host.Limit > 0 { requestsToReplay += host.Limit } else { From bba9637eed34253d72251fa9de8fbd0d6fd6b95b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C5=82awosz?= Date: Tue, 8 Oct 2013 19:54:37 +0200 Subject: [PATCH 31/31] remove comment --- integration_test.go | 2 -- 1 file changed, 2 deletions(-) diff --git a/integration_test.go b/integration_test.go index 5916e18..718ca46 100644 --- a/integration_test.go +++ b/integration_test.go @@ -205,8 +205,6 @@ func (e *Env) startFileListener() (p int) { go e.startHTTP(p, http.HandlerFunc(e.ListenHandler)) go e.startHTTP(p+2, http.HandlerFunc(e.ReplayHandler)) go e.startFileUsingListener(p, p+1) - // we will replay after listener finishes capturing - // go e.startFileUsingReplay(p+1, p+2) // Time to start http and gor instances time.Sleep(time.Millisecond * 100)