diff --git a/raw_socket_listener/listener.go b/raw_socket_listener/listener.go index cecf979..1b1c8c0 100644 --- a/raw_socket_listener/listener.go +++ b/raw_socket_listener/listener.go @@ -33,6 +33,12 @@ import ( var _ = fmt.Println +type Packet struct { + srcIP []byte + data []byte + timestamp time.Time +} + // Listener handle traffic capture type Listener struct { mu sync.Mutex @@ -53,7 +59,7 @@ type Listener struct { respWithoutReq map[uint32]tcpID // Messages ready to be send to client - packetsChan chan []byte + packetsChan chan *Packet // Messages ready to be send to client messagesChan chan *TCPMessage @@ -88,7 +94,7 @@ const ( func NewListener(addr string, port string, engine int, trackResponse bool, expire time.Duration) (l *Listener) { l = &Listener{} - l.packetsChan = make(chan []byte, 10000) + l.packetsChan = make(chan *Packet, 10000) l.messagesChan = make(chan *TCPMessage, 10000) l.quit = make(chan bool) l.readyCh = make(chan bool, 1) @@ -137,9 +143,9 @@ func (t *Listener) listen() { t.conn.Close() } return - case data := <-t.packetsChan: - packet := ParseTCPPacket(data[:16], data[16:]) - t.processTCPPacket(packet) + case packet := <-t.packetsChan: + tcpPacket := ParseTCPPacket(packet.srcIP, packet.data, packet.timestamp) + t.processTCPPacket(tcpPacket) case <-gcTicker: now := time.Now() @@ -522,11 +528,13 @@ func (t *Listener) readPcap() { } } - newBuf := make([]byte, len(data)+16) - copy(newBuf[:16], srcIP) - copy(newBuf[16:], data) + packetSrcIP := make([]byte, 16) + packetData := make([]byte, len(data)) - t.packetsChan <- newBuf + copy(packetSrcIP, srcIP) + copy(packetData, data) + + t.packetsChan <- t.buildPacket(srcIP, data, packet.Metadata().Timestamp) } } }(d) @@ -589,11 +597,7 @@ func (t *Listener) readPcapFile() { continue } - newBuf := make([]byte, len(data)+16) - copy(newBuf[:16], addr) - copy(newBuf[16:], data) - - t.packetsChan <- newBuf + t.packetsChan <- t.buildPacket(addr, data, packet.Metadata().Timestamp) } } } @@ -626,16 +630,26 @@ func (t *Listener) readRAWSocket() { if n > 0 { if t.isValidPacket(buf[:n]) { - newBuf := make([]byte, n+16) - copy(newBuf[16:], buf[:n]) - copy(newBuf[:16], []byte(addr.(*net.IPAddr).IP)) - - t.packetsChan <- newBuf + t.packetsChan <- t.buildPacket([]byte(addr.(*net.IPAddr).IP), buf[:n], time.Now()) } } } } +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, + timestamp:timestamp, + } +} + func (t *Listener) isValidPacket(buf []byte) bool { // To avoid full packet parsing every time, we manually parsing values needed for packet filtering // http://en.wikipedia.org/wiki/Transmission_Control_Protocol @@ -718,7 +732,7 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { message, ok := t.messages[packet.ID] if !ok { - message = NewTCPMessage(packet.Seq, packet.Ack, isIncoming) + message = NewTCPMessage(packet.Seq, packet.Ack, isIncoming, packet.timestamp) t.messages[packet.ID] = message if !isIncoming { diff --git a/raw_socket_listener/listener_test.go b/raw_socket_listener/listener_test.go index 72e01de..8ade26e 100644 --- a/raw_socket_listener/listener_test.go +++ b/raw_socket_listener/listener_test.go @@ -15,10 +15,10 @@ func TestRawListenerInput(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")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) respAck := reqPacket.Seq + uint32(len(reqPacket.Data)) - respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) listener.packetsChan <- reqPacket.Dump() listener.packetsChan <- respPacket.Dump() @@ -52,11 +52,11 @@ func TestRawListenerInputResponseByClose(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")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) respAck := reqPacket.Seq + uint32(len(reqPacket.Data)) - respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\nConnection: close\r\n\r\nasd")) - finPacket := buildPacket(false, respAck, reqPacket.Seq+2, []byte("")) + respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\nConnection: close\r\n\r\nasd"), time.Now()) + finPacket := buildPacket(false, respAck, reqPacket.Seq+2, []byte(""), time.Now()) finPacket.IsFIN = true listener.packetsChan <- reqPacket.Dump() @@ -92,7 +92,7 @@ func TestRawListenerInputWithoutResponse(t *testing.T) { 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")) + reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) listener.packetsChan <- reqPacket.Dump() @@ -114,8 +114,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")) - respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + 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()) // If response packet comes before request listener.packetsChan <- respPacket.Dump() @@ -152,15 +152,15 @@ 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")) + 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")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + 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")) + 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")) + 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") @@ -172,15 +172,15 @@ 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")) + 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")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + 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")) + 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")) + 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") @@ -191,15 +191,15 @@ 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")) + 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")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+1, []byte("b")) + 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")) + 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")) + 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") @@ -352,16 +352,16 @@ 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")) + 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")) - reqPacket3 := buildPacket(true, 2, reqPacket2.Seq+uint32(len(reqPacket2.Data)), []byte("1\r\nb\r\n")) - reqPacket4 := buildPacket(true, 2, reqPacket3.Seq+uint32(len(reqPacket3.Data)), []byte("0\r\n\r\n")) + 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()) - respPacket1 := buildPacket(false, 10, 3, []byte("HTTP/1.1 100 Continue\r\n\r\n")) + 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, reqPacket4.Seq+5 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket2 := buildPacket(false, reqPacket4.Seq+5 /* len of data */, 2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) // Should re-construct message from all possible combinations for i := 0; i < 6*5*4*3*2*1; i++ { @@ -381,13 +381,13 @@ 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")) + reqPacket1 := buildPacket(true, ack, seq, []byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n"), time.Now()) // 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")) + reqPacket2 := buildPacket(true, ack, seq+47, []byte("1\r\na\r\n"), time.Now()) + reqPacket3 := buildPacket(true, ack, reqPacket2.Seq+5, []byte("1\r\nb\r\n"), time.Now()) + reqPacket4 := buildPacket(true, ack, reqPacket3.Seq+5, []byte("0\r\n\r\n"), time.Now()) - respPacket := buildPacket(false, reqPacket4.Seq+5 /* len of data */, ack, []byte("HTTP/1.1 200 OK\r\n\r\n")) + respPacket := buildPacket(false, reqPacket4.Seq+5 /* len of data */, ack, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) return []*TCPPacket{ reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket, @@ -409,8 +409,8 @@ func postMessage() []*TCPPacket { } return []*TCPPacket{ - buildPacket(true, ack, seq, data), - buildPacket(false, seq+uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n\r\n")), + buildPacket(true, ack, seq, data, time.Now()), + buildPacket(false, seq+uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()), } } @@ -420,8 +420,8 @@ func getMessage() []*TCPPacket { 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\r\n")), + buildPacket(true, ack, seq, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()), + buildPacket(false, seq+18, seq2, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()), } } diff --git a/raw_socket_listener/tcp_message.go b/raw_socket_listener/tcp_message.go index 5cb566e..5319594 100644 --- a/raw_socket_listener/tcp_message.go +++ b/raw_socket_listener/tcp_message.go @@ -48,9 +48,8 @@ type TCPMessage struct { } // NewTCPMessage pointer created from a Acknowledgment number and a channel of messages readuy to be deleted -func NewTCPMessage(Seq, Ack uint32, IsIncoming bool) (msg *TCPMessage) { - msg = &TCPMessage{Seq: Seq, Ack: Ack, IsIncoming: IsIncoming} - msg.Start = time.Now() +func NewTCPMessage(Seq, Ack uint32, IsIncoming bool, timestamp time.Time) (msg *TCPMessage) { + msg = &TCPMessage{Seq: Seq, Ack: Ack, IsIncoming: IsIncoming, Start: timestamp} return } @@ -138,6 +137,10 @@ func (t *TCPMessage) AddPacket(packet *TCPPacket) { if packet.OrigAck != 0 { t.DataAck = packet.OrigAck } + + if packet.timestamp.Before(t.Start) { + t.Start = packet.timestamp + } } t.checkSeqIntegrity() diff --git a/raw_socket_listener/tcp_message_test.go b/raw_socket_listener/tcp_message_test.go index 781a824..1ce0d19 100644 --- a/raw_socket_listener/tcp_message_test.go +++ b/raw_socket_listener/tcp_message_test.go @@ -5,9 +5,10 @@ import ( "encoding/binary" _ "log" "testing" + "time" ) -func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPacket) { +func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte, timestamp time.Time) (packet *TCPPacket) { var srcPort, destPort uint16 // For tests `listening` port is 0 @@ -25,7 +26,7 @@ func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPack buf[12] = 64 buf = append(buf, Data...) - packet = ParseTCPPacket([]byte("123"), buf) + packet = ParseTCPPacket([]byte("123"), buf, timestamp) return packet } @@ -36,31 +37,31 @@ func buildMessage(p *TCPPacket) *TCPMessage { isIncoming = true } - m := NewTCPMessage(p.Seq, p.Ack, isIncoming) + m := NewTCPMessage(p.Seq, p.Ack, isIncoming, p.timestamp) m.AddPacket(p) return m } func TestTCPMessagePacketsOrder(t *testing.T) { - msg := buildMessage(buildPacket(true, 1, 1, []byte("a"))) - msg.AddPacket(buildPacket(true, 1, 2, []byte("b"))) + msg := buildMessage(buildPacket(true, 1, 1, []byte("a"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 2, []byte("b"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("ab")) { t.Error("Should contatenate packets in right order") } // When first packet have wrong order (Seq) - msg = buildMessage(buildPacket(true, 1, 2, []byte("b"))) - msg.AddPacket(buildPacket(true, 1, 1, []byte("a"))) + msg = buildMessage(buildPacket(true, 1, 2, []byte("b"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("a"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("ab")) { t.Error("Should contatenate packets in right order") } // Should ignore packets with same sequence - msg = buildMessage(buildPacket(true, 1, 1, []byte("a"))) - msg.AddPacket(buildPacket(true, 1, 1, []byte("a"))) + msg = buildMessage(buildPacket(true, 1, 1, []byte("a"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("a"), time.Now())) if !bytes.Equal(msg.Bytes(), []byte("a")) { t.Error("Should ignore packet with same Seq") @@ -68,8 +69,8 @@ func TestTCPMessagePacketsOrder(t *testing.T) { } func TestTCPMessageSize(t *testing.T) { - msg := buildMessage(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"))) - msg.AddPacket(buildPacket(true, 1, 2, []byte("b"))) + msg := buildMessage(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 2, []byte("b"), time.Now())) if msg.BodySize() != 2 { t.Error("Should count only body", msg.BodySize()) @@ -110,7 +111,7 @@ func TestTCPMessageIsComplete(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if tc.assocMessage { msg.AssocMessage = &TCPMessage{} } @@ -123,9 +124,9 @@ func TestTCPMessageIsComplete(t *testing.T) { } func TestTCPMessageIsSeqMissing(t *testing.T) { - p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n")) - p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n")) - p3 := buildPacket(false, 1, p2.Seq+uint32(len(p2.Data)), []byte("a")) + p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n"), time.Now()) + p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n"), time.Now()) + p3 := buildPacket(false, 1, p2.Seq+uint32(len(p2.Data)), []byte("a"), time.Now()) msg := buildMessage(p1) if msg.seqMissing { @@ -144,8 +145,8 @@ func TestTCPMessageIsSeqMissing(t *testing.T) { } func TestTCPMessageIsHeadersReceived(t *testing.T) { - p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n\r\n")) - p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n")) + p1 := buildPacket(false, 1, 1, []byte("HTTP/1.1 200 OK\r\n\r\n"), time.Now()) + p2 := buildPacket(false, 1, p1.Seq+uint32(len(p1.Data)), []byte("Content-Length: 10\r\n\r\n"), time.Now()) msg := buildMessage(p1) if msg.headerPacket == -1 { @@ -157,7 +158,7 @@ func TestTCPMessageIsHeadersReceived(t *testing.T) { t.Error("Should found double new line: headers received") } - msg = buildMessage(buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\nContent-Length: 1\r\n"))) + msg = buildMessage(buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\nContent-Length: 1\r\n"), time.Now())) if msg.headerPacket != -1 { t.Error("Should not find headers end") } @@ -183,7 +184,7 @@ func TestTCPMessageMethodType(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if msg.methodType != tc.expectedMethodType { t.Errorf("Expected %d, got %d", tc.expectedMethodType, msg.methodType) @@ -208,7 +209,7 @@ func TestTCPMessageBodyType(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payload), time.Now())) if msg.bodyType != tc.expectedBodyType { t.Errorf("Expected %d, got %d", tc.expectedBodyType, msg.bodyType) @@ -229,12 +230,12 @@ func TestTCPMessageBodySize(t *testing.T) { } for _, tc := range testCases { - msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payloads[0]))) + msg := buildMessage(buildPacket(tc.direction, 1, 1, []byte(tc.payloads[0]), time.Now())) if len(tc.payloads) > 1 { for _, p := range tc.payloads[1:] { seq := uint32(1 + msg.Size()) - msg.AddPacket(buildPacket(tc.direction, 1, seq, []byte(p))) + msg.AddPacket(buildPacket(tc.direction, 1, seq, []byte(p), time.Now())) } } @@ -243,3 +244,15 @@ func TestTCPMessageBodySize(t *testing.T) { } } } + +func TestTcpMessageStart(t *testing.T) { + start := time.Now().Add(-1 * time.Second) + + msg := buildMessage(buildPacket(true, 1, 2, []byte("b"), time.Now())) + msg.AddPacket(buildPacket(true, 1, 1, []byte("POST / HTTP/1.1\r\nContent-Length: 2\r\n\r\na"), start)) + + if msg.Start != start { + t.Error("Message timestamp should be equal to the lowest related packet timestamp", start, msg.Start) + } +} + diff --git a/raw_socket_listener/tcp_packet.go b/raw_socket_listener/tcp_packet.go index 0e3e832..1dafcbe 100644 --- a/raw_socket_listener/tcp_packet.go +++ b/raw_socket_listener/tcp_packet.go @@ -5,6 +5,7 @@ import ( "log" "strconv" "strings" + "time" ) var _ = log.Println @@ -38,14 +39,16 @@ type TCPPacket struct { Raw []byte Data []byte Addr []byte + timestamp time.Time ID tcpID } // ParseTCPPacket takes address and tcp payload and returns parsed TCPPacket -func ParseTCPPacket(addr []byte, data []byte) (p *TCPPacket) { +func ParseTCPPacket(addr []byte, data []byte, timestamp time.Time) (p *TCPPacket) { p = &TCPPacket{Raw: data} p.ParseBasic() p.Addr = addr + p.timestamp = timestamp p.GenID() return @@ -79,27 +82,33 @@ func (t *TCPPacket) ParseBasic() { t.Data = t.Raw[t.DataOffset*4:] } -func (t *TCPPacket) Dump() []byte { - buf := make([]byte, len(t.Data)+16+16) - copy(buf[:16], t.Addr) +func (t *TCPPacket) Dump() *Packet { - tcpBuf := buf[16:] + packetSrcIP := make([]byte, 16) + packetData := make([]byte, len(t.Data) + 16) - binary.BigEndian.PutUint16(tcpBuf[2:4], t.DestPort) - binary.BigEndian.PutUint16(tcpBuf[0:2], t.SrcPort) + copy(packetSrcIP, t.Addr) - binary.BigEndian.PutUint32(tcpBuf[4:8], t.Seq) - binary.BigEndian.PutUint32(tcpBuf[8:12], t.Ack) + binary.BigEndian.PutUint16(packetData[0:2], t.SrcPort) + binary.BigEndian.PutUint16(packetData[2:4], t.DestPort) - tcpBuf[12] = 64 + binary.BigEndian.PutUint32(packetData[4:8], t.Seq) + binary.BigEndian.PutUint32(packetData[8:12], t.Ack) + + packetData[12] = 64 if t.IsFIN { - tcpBuf[13] = tcpBuf[13] | 0x01 + packetData[13] = packetData[13] | 0x01 } - copy(tcpBuf[16:], t.Data) + copy(packetData[16:], t.Data) + + return &Packet{ + srcIP: packetSrcIP, + data:packetData, + timestamp:t.timestamp, + } - return buf } // String output for a TCP Packet