From d62eccb655f048ac8434094e2b41ba2b20564b3e Mon Sep 17 00:00:00 2001 From: Leonid Bugaev Date: Tue, 21 Jul 2015 10:06:10 +0500 Subject: [PATCH] Make raw message expiration configurable to improve tests time --- input_raw.go | 7 +++++-- input_raw_test.go | 10 ++++++---- middleware_test.go | 4 ++-- raw_socket_listener/listener.go | 13 +++++++++++-- raw_socket_listener/tcp_message.go | 13 ++++++------- 5 files changed, 30 insertions(+), 17 deletions(-) diff --git a/input_raw.go b/input_raw.go index c4f3c34..613cbbe 100644 --- a/input_raw.go +++ b/input_raw.go @@ -5,19 +5,22 @@ import ( "log" "net" "strings" + "time" ) // RAWInput used for intercepting traffic for given address type RAWInput struct { data chan []byte address string + expire time.Duration } // NewRAWInput constructor for RAWInput. Accepts address with port as argument. -func NewRAWInput(address string) (i *RAWInput) { +func NewRAWInput(address string, expire time.Duration) (i *RAWInput) { i = new(RAWInput) i.data = make(chan []byte) i.address = address + i.expire = expire go i.listen(address) @@ -42,7 +45,7 @@ func (i *RAWInput) listen(address string) { log.Fatal("input-raw: error while parsing address", err) } - listener := raw.NewListener(host, port) + listener := raw.NewListener(host, port, i.expire) for { // Receiving TCPMessage object diff --git a/input_raw_test.go b/input_raw_test.go index 0f7b0c3..a145845 100644 --- a/input_raw_test.go +++ b/input_raw_test.go @@ -15,6 +15,8 @@ import ( "time" ) +const testRawExpire = time.Millisecond * 100 + func TestRAWInput(t *testing.T) { wg := new(sync.WaitGroup) @@ -22,7 +24,7 @@ func TestRAWInput(t *testing.T) { listener := startHTTP(func(w http.ResponseWriter, req *http.Request) {}) - input := NewRAWInput(listener.Addr().String()) + input := NewRAWInput(listener.Addr().String(), testRawExpire) output := NewTestOutput(func(data []byte) { wg.Done() }) @@ -64,7 +66,7 @@ func TestInputRAW100Expect(t *testing.T) { originAddr := strings.Replace(origin.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr) + input := NewRAWInput(originAddr, testRawExpire) // We will use it to get content of raw HTTP request testOutput := NewTestOutput(func(data []byte) { @@ -121,7 +123,7 @@ func TestInputRAWChunkedEncoding(t *testing.T) { originAddr := strings.Replace(origin.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr) + input := NewRAWInput(originAddr, testRawExpire) listener := startHTTP(func(w http.ResponseWriter, req *http.Request) { defer req.Body.Close() @@ -179,7 +181,7 @@ func TestInputRAWLargePayload(t *testing.T) { })) originAddr := strings.Replace(origin.Listener.Addr().String(), "[::]", "127.0.0.1", -1) - input := NewRAWInput(originAddr) + input := NewRAWInput(originAddr, testRawExpire) replay := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { req.Body = http.MaxBytesReader(w, req.Body, 1*1024*1024) diff --git a/middleware_test.go b/middleware_test.go index cfea8b4..cb7319a 100644 --- a/middleware_test.go +++ b/middleware_test.go @@ -105,7 +105,7 @@ func TestEchoMiddleware(t *testing.T) { quit := make(chan int) // Catch traffic from one service - input := NewRAWInput(from.Listener.Addr().String()) + input := NewRAWInput(from.Listener.Addr().String(), testRawExpire) // And redirect to another output := NewHTTPOutput(to.URL, &HTTPOutputConfig{}) @@ -156,7 +156,7 @@ func TestTokenMiddleware(t *testing.T) { quit := make(chan int) // Catch traffic from one service - input := NewRAWInput(from) + input := NewRAWInput(from, testRawExpire) // And redirect to another output := NewHTTPOutput(to, &HTTPOutputConfig{}) diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index f82ae6f..53857ac 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -18,6 +18,7 @@ import ( "log" "net" "strconv" + "time" ) // Listener handle traffic capture @@ -42,10 +43,12 @@ type Listener struct { addr string // IP to listen port int // Port to listen + + messageExpire time.Duration } // NewListener creates and initializes new Listener object -func NewListener(addr string, port string) (rawListener *Listener) { +func NewListener(addr string, port string, expire time.Duration) (rawListener *Listener) { rawListener = &Listener{} rawListener.packetsChan = make(chan *TCPPacket, 10000) @@ -59,6 +62,12 @@ func NewListener(addr string, port string) (rawListener *Listener) { rawListener.addr = addr rawListener.port, _ = strconv.Atoi(port) + if expire.Nanoseconds() == 0 { + expire = 2000 * time.Millisecond + } + + rawListener.messageExpire = expire + go rawListener.listen() go rawListener.readRAWSocket() @@ -157,7 +166,7 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { if !ok { // We sending messageDelChan channel, so message object can communicate with Listener and notify it if message completed - message = NewTCPMessage(mID, t.messageDelChan, packet.Ack) + message = NewTCPMessage(mID, t.messageDelChan, packet.Ack, &t.messageExpire) t.messages[mID] = message } diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index ddfb0f2..e6e31bd 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -6,9 +6,6 @@ import ( "time" ) -// MsgExpire specify period that message should wait before it considered as finished -const MsgExpire = 2000 * time.Millisecond - // TCPMessage ensure that all TCP packets for given request is received, and processed in right sequence // Its needed because all TCP message can be fragmented or re-transmitted // @@ -25,17 +22,19 @@ type TCPMessage struct { packetsChan chan *TCPPacket delChan chan *TCPMessage + + expire *time.Duration } // NewTCPMessage pointer created from a Acknowledgment number and a channel of messages readuy to be deleted -func NewTCPMessage(ID string, delChan chan *TCPMessage, Ack uint32) (msg *TCPMessage) { - msg = &TCPMessage{ID: ID, Ack: Ack} +func NewTCPMessage(ID string, delChan chan *TCPMessage, Ack uint32, expire *time.Duration) (msg *TCPMessage) { + msg = &TCPMessage{ID: ID, Ack: Ack, expire: expire} msg.packetsChan = make(chan *TCPPacket) msg.delChan = delChan // used for notifying that message completed or expired // Every time we receive packet we reset this timer - msg.timer = time.AfterFunc(MsgExpire, msg.Timeout) + msg.timer = time.AfterFunc(*msg.expire, msg.Timeout) go msg.listen() @@ -103,5 +102,5 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { } // Reset message timeout timer - t.timer.Reset(MsgExpire) + t.timer.Reset(*t.expire) }