mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Merge pull request #395 from ylegat/feature/use-pcap-timestamp
Resolve buger/gor#392 : use pcap timestamp
This commit is contained in:
@@ -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,7 @@ func (t *Listener) readPcap() {
|
||||
}
|
||||
}
|
||||
|
||||
newBuf := make([]byte, len(data)+16)
|
||||
copy(newBuf[:16], srcIP)
|
||||
copy(newBuf[16:], data)
|
||||
|
||||
t.packetsChan <- newBuf
|
||||
t.packetsChan <- t.buildPacket(srcIP, data, packet.Metadata().Timestamp)
|
||||
}
|
||||
}
|
||||
}(d)
|
||||
@@ -589,11 +591,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 +624,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 +726,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 {
|
||||
|
||||
@@ -15,13 +15,13 @@ 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()
|
||||
listener.packetsChan <- reqPacket.dump()
|
||||
listener.packetsChan <- respPacket.dump()
|
||||
|
||||
select {
|
||||
case req = <-listener.messagesChan:
|
||||
@@ -52,16 +52,16 @@ 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()
|
||||
listener.packetsChan <- respPacket.Dump()
|
||||
listener.packetsChan <- finPacket.Dump()
|
||||
listener.packetsChan <- reqPacket.dump()
|
||||
listener.packetsChan <- respPacket.dump()
|
||||
listener.packetsChan <- finPacket.dump()
|
||||
|
||||
select {
|
||||
case req = <-listener.messagesChan:
|
||||
@@ -92,9 +92,9 @@ 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()
|
||||
listener.packetsChan <- reqPacket.dump()
|
||||
|
||||
select {
|
||||
case req = <-listener.messagesChan:
|
||||
@@ -114,12 +114,12 @@ 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()
|
||||
listener.packetsChan <- reqPacket.Dump()
|
||||
listener.packetsChan <- respPacket.dump()
|
||||
listener.packetsChan <- reqPacket.dump()
|
||||
|
||||
select {
|
||||
case req = <-listener.messagesChan:
|
||||
@@ -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")
|
||||
|
||||
@@ -209,7 +209,7 @@ func TestAlt100ContinueHeaderOrder(t *testing.T) {
|
||||
func testRawListener100Continue(t *testing.T, listener *Listener, result []byte, packets ...*TCPPacket) {
|
||||
var req, resp *TCPMessage
|
||||
for _, p := range packets {
|
||||
listener.packetsChan <- p.Dump()
|
||||
listener.packetsChan <- p.dump()
|
||||
}
|
||||
|
||||
select {
|
||||
@@ -249,7 +249,7 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket
|
||||
var r, req, resp *TCPMessage
|
||||
|
||||
for _, p := range packets {
|
||||
listener.packetsChan <- p.Dump()
|
||||
listener.packetsChan <- p.dump()
|
||||
}
|
||||
|
||||
select {
|
||||
@@ -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()),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -452,7 +452,7 @@ func TestRawListenerBench(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
l.packetsChan <- p.Dump()
|
||||
l.packetsChan <- p.dump()
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user