mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Add dest port to id
This commit is contained in:
@@ -475,7 +475,7 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) {
|
||||
// If message contains only single packet immediately dispatch it
|
||||
if message.IsFinished() {
|
||||
if isIncoming {
|
||||
if resp, ok := t.messages[message.ResponseID()]; ok {
|
||||
if resp, ok := t.messages[message.ResponseID]; ok {
|
||||
t.dispatchMessage(message)
|
||||
if resp.IsFinished() {
|
||||
t.dispatchMessage(resp)
|
||||
|
||||
@@ -15,10 +15,10 @@ func TestRawListenerInput(t *testing.T) {
|
||||
listener := NewListener("", "0", EnginePcap, 10*time.Millisecond)
|
||||
defer listener.Close()
|
||||
|
||||
reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1"))
|
||||
reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"))
|
||||
|
||||
respAck := reqPacket.Seq + uint32(len(reqPacket.Data))
|
||||
respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK"))
|
||||
respPacket := buildPacket(false, respAck, reqPacket.Seq+1, []byte("HTTP/1.1 200 OK\r\n\r\n"))
|
||||
|
||||
listener.processTCPPacket(reqPacket)
|
||||
listener.processTCPPacket(respPacket)
|
||||
@@ -52,8 +52,8 @@ func TestRawListenerResponse(t *testing.T) {
|
||||
listener := NewListener("", "0", EnginePcap, 10*time.Millisecond)
|
||||
defer listener.Close()
|
||||
|
||||
reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1"))
|
||||
respPacket := buildPacket(false, 1+uint32(len(reqPacket.Data)), 2, []byte("HTTP/1.1 200 OK"))
|
||||
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"))
|
||||
|
||||
// If response packet comes before request
|
||||
listener.processTCPPacket(respPacket)
|
||||
|
||||
@@ -11,6 +11,8 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
var _ = log.Println
|
||||
|
||||
// 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
|
||||
//
|
||||
@@ -24,6 +26,7 @@ type TCPMessage struct {
|
||||
RequestStart time.Time
|
||||
RequestAck uint32
|
||||
RequestID tcpID
|
||||
ResponseID tcpID
|
||||
Start time.Time
|
||||
End time.Time
|
||||
IsIncoming bool
|
||||
@@ -190,22 +193,21 @@ func (t *TCPMessage) UUID() []byte {
|
||||
// UpdateResponseAck should be called after packet is added
|
||||
func (t *TCPMessage) UpdateResponseAck() uint32 {
|
||||
lastPacket := t.packets[len(t.packets)-1]
|
||||
t.ResponseAck = lastPacket.Seq + uint32(len(lastPacket.Data))
|
||||
respAck := lastPacket.Seq + uint32(len(lastPacket.Data))
|
||||
|
||||
if t.ResponseAck != respAck {
|
||||
t.ResponseAck = lastPacket.Seq + uint32(len(lastPacket.Data))
|
||||
|
||||
// We swappwed src and dst port
|
||||
copy(t.ResponseID[:4], lastPacket.Addr)
|
||||
copy(t.ResponseID[4:], lastPacket.Raw[2:4]) // Src port
|
||||
copy(t.ResponseID[6:], lastPacket.Raw[0:2]) // Dest port
|
||||
binary.BigEndian.PutUint32(t.ResponseID[8:12], t.ResponseAck)
|
||||
}
|
||||
|
||||
return t.ResponseAck
|
||||
}
|
||||
|
||||
func (t *TCPMessage) ID() tcpID {
|
||||
return t.packets[0].ID
|
||||
}
|
||||
|
||||
// ResponseID generate message ID for request response
|
||||
func (t *TCPMessage) ResponseID() tcpID {
|
||||
var id tcpID
|
||||
p := t.packets[0]
|
||||
|
||||
copy(id[:4], p.Addr)
|
||||
copy(id[4:], p.Data[2:4]) // Dest port
|
||||
binary.BigEndian.PutUint32(id[6:10], t.ResponseAck)
|
||||
|
||||
return id
|
||||
}
|
||||
}
|
||||
@@ -4,23 +4,29 @@ import (
|
||||
"bytes"
|
||||
_ "log"
|
||||
"testing"
|
||||
"encoding/binary"
|
||||
)
|
||||
|
||||
func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPacket) {
|
||||
packet = &TCPPacket{
|
||||
Addr: []byte(""),
|
||||
Ack: Ack,
|
||||
Seq: Seq,
|
||||
Data: Data,
|
||||
}
|
||||
var srcPort, destPort uint16
|
||||
|
||||
// For tests `listening` port is 0
|
||||
if isIncoming {
|
||||
packet.SrcPort = 1
|
||||
srcPort = 1
|
||||
} else {
|
||||
packet.DestPort = 1
|
||||
destPort = 1
|
||||
}
|
||||
|
||||
buf := make([]byte, 16)
|
||||
binary.BigEndian.PutUint16(buf[2:4], destPort)
|
||||
binary.BigEndian.PutUint16(buf[0:2], srcPort)
|
||||
binary.BigEndian.PutUint32(buf[4:8], Seq)
|
||||
binary.BigEndian.PutUint32(buf[8:12], Ack)
|
||||
buf[12] = 64
|
||||
buf = append(buf, Data...)
|
||||
|
||||
packet = ParseTCPPacket([]byte(""), buf)
|
||||
|
||||
return packet
|
||||
}
|
||||
|
||||
|
||||
@@ -19,7 +19,7 @@ const (
|
||||
fNS
|
||||
)
|
||||
|
||||
type tcpID [10]byte
|
||||
type tcpID [12]byte
|
||||
|
||||
// TCPPacket provides tcp packet parser
|
||||
// Packet structure: http://en.wikipedia.org/wiki/Transmission_Control_Protocol
|
||||
@@ -43,8 +43,9 @@ func ParseTCPPacket(addr []byte, data []byte) (p *TCPPacket) {
|
||||
p.Addr = addr
|
||||
|
||||
copy(p.ID[:4], addr)
|
||||
copy(p.ID[4:], p.Raw[2:4]) // Dest port
|
||||
copy(p.ID[6:], p.Raw[8:12]) // Ack
|
||||
copy(p.ID[4:], p.Raw[0:2]) // Src port
|
||||
copy(p.ID[6:], p.Raw[2:4]) // Dest port
|
||||
copy(p.ID[8:], p.Raw[8:12]) // Ack
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user