From 0b437716bedce0d91726f8a0ef14ff230a9fe62d Mon Sep 17 00:00:00 2001 From: Urban Ishimwe Date: Sat, 6 Jun 2020 09:51:19 +0200 Subject: [PATCH 1/3] make raw socket usable, minor optimazation and fixed tests. - ListenPacket can now capture request(non-multicast) --- Makefile | 4 +- {raw_socket_listener => capture}/listener.go | 103 +++++++++--------- .../listener_test.go | 55 +++++----- .../tcp_message.go | 5 +- .../tcp_message_test.go | 2 +- .../tcp_packet.go | 6 +- input_raw.go | 12 +- 7 files changed, 85 insertions(+), 102 deletions(-) rename {raw_socket_listener => capture}/listener.go (93%) rename {raw_socket_listener => capture}/listener_test.go (94%) rename {raw_socket_listener => capture}/tcp_message.go (99%) rename {raw_socket_listener => capture}/tcp_message_test.go (99%) rename {raw_socket_listener => capture}/tcp_packet.go (98%) diff --git a/Makefile b/Makefile index a663ea9..64facfd 100644 --- a/Makefile +++ b/Makefile @@ -63,8 +63,8 @@ bench: $(RUN) go test $(LDFLAGS) -v -run NOT_EXISTING -bench $(BENCHMARK) -benchtime 5s profile_test: - $(RUN) go test $(LDFLAGS) -run $(TEST) ./raw_socket_listener/. $(ARGS) -memprofile mem.mprof -cpuprofile cpu.out - $(RUN) go test $(LDFLAGS) -run $(TEST) ./raw_socket_listener/. $(ARGS) -c + $(RUN) go test $(LDFLAGS) -run $(TEST) ./capture/. $(ARGS) -memprofile mem.mprof -cpuprofile cpu.out + $(RUN) go test $(LDFLAGS) -run $(TEST) ./capture/. $(ARGS) -c # Used mainly for debugging, because docker container do not have access to parent machine ports run: diff --git a/raw_socket_listener/listener.go b/capture/listener.go similarity index 93% rename from raw_socket_listener/listener.go rename to capture/listener.go index d3f42e6..d3cc70b 100644 --- a/raw_socket_listener/listener.go +++ b/capture/listener.go @@ -1,21 +1,15 @@ /* -Package rawSocket provides traffic sniffier using RAW sockets. - -Capture traffic from socket using RAW_SOCKET's -http://en.wikipedia.org/wiki/Raw_socket - -RAW_SOCKET allow you listen for traffic on any port (e.g. sniffing) because they operate on IP level. - +Package capture provides traffic sniffier using RAW sockets. +Capture traffic from socket using RAW_SOCKET's http://en.wikipedia.org/wiki/Raw_socket +RAW_SOCKET allows you to listen for traffic from any port (e.g. sniffing) because they operate on IP level. Ports is TCP feature, same as flow control, reliable transmission and etc. - This package implements own TCP layer: TCP packets is parsed using tcp_packet.go, and flow control is managed by tcp_message.go */ -package rawSocket +package capture import ( "bytes" "encoding/binary" - "fmt" "io" "log" "net" @@ -33,8 +27,6 @@ import ( "github.com/google/gopacket/pcap" ) -var _ = fmt.Println - type packet struct { srcIP []byte data []byte @@ -43,7 +35,7 @@ type packet struct { // Listener handle traffic capture type Listener struct { - mu sync.Mutex + sync.Mutex // buffer of TCPMessages waiting to be send // ID -> TCPMessage messages map[tcpID]*TCPMessage @@ -83,7 +75,7 @@ type Listener struct { pcapHandles []*pcap.Handle quit chan bool - readyCh chan bool + readyCh bool } type request struct { @@ -106,7 +98,6 @@ func NewListener(addr string, port string, engine int, trackResponse bool, expir l.packetsChan = make(chan *packet, 10000) l.messagesChan = make(chan *TCPMessage, 10000) l.quit = make(chan bool) - l.readyCh = make(chan bool, 1) l.messages = make(map[tcpID]*TCPMessage) l.ackAliases = make(map[uint32]uint32) @@ -139,6 +130,8 @@ func NewListener(addr string, port string, engine int, trackResponse bool, expir go l.readPcap() case EnginePcapFile: go l.readPcapFile() + case EngineRawSocket: + go l.readRAWSocket() default: log.Fatal("Unknown traffic interception engine:", engine) } @@ -340,6 +333,7 @@ func (t *Listener) readPcap() { go func(device pcap.Interface) { inactive, err := pcap.NewInactiveHandle(device.Name) if err != nil { + inactive.CleanUp() log.Println("Pcap Error while opening device", device.Name, err) wg.Done() return @@ -359,7 +353,7 @@ func (t *Listener) readPcap() { } else { inactive.SetSnapLen(65536) } - + inactive.SetSnapLen(65536) inactive.SetTimeout(t.messageExpire) inactive.SetPromisc(true) inactive.SetImmediateMode(t.immediateMode) @@ -379,7 +373,7 @@ func (t *Listener) readPcap() { defer handle.Close() - t.mu.Lock() + t.Lock() t.pcapHandles = append(t.pcapHandles, handle) var bpfDstHost, bpfSrcHost string @@ -425,7 +419,7 @@ func (t *Listener) readPcap() { return } } - t.mu.Unlock() + t.Unlock() var decoder gopacket.Decoder @@ -526,7 +520,7 @@ func (t *Listener) readPcap() { continue } - dataOffset := (data[12] & 0xF0) >> 4 + dataOffset := data[12] >> 4 isFIN := data[13]&0x01 != 0 // We need only packets with data inside @@ -586,7 +580,9 @@ func (t *Listener) readPcap() { } wg.Wait() - t.readyCh <- true + t.Lock() + t.readyCh = true + t.Unlock() } func (t *Listener) readPcapFile() { @@ -600,7 +596,9 @@ func (t *Listener) readPcapFile() { } } - t.readyCh <- true + t.Lock() + t.readyCh = true + t.Unlock() packetSource := gopacket.NewPacketSource(handle, handle.LinkType()) for { @@ -640,7 +638,7 @@ func (t *Listener) readPcapFile() { continue } - dataOffset := (data[12] & 0xF0) >> 4 + dataOffset := data[12] >> 4 isFIN := data[13]&0x01 != 0 // We need only packets with data inside @@ -657,32 +655,45 @@ func (t *Listener) readPcapFile() { func (t *Listener) readRAWSocket() { conn, e := net.ListenPacket("ip:tcp", t.addr) t.conn = conn - if e != nil { log.Fatal(e) } - defer t.conn.Close() - - buf := make([]byte, 64*1024) // 64kb - - t.readyCh <- true - + type RSPacket struct { + buf []byte + addr net.Addr + err error + n int + } + var bufChan = make(chan *RSPacket, 1000) + t.Lock() + t.readyCh = true + t.Unlock() + go func() { + buffer := 64 * 1024 + if t.bufferSize > int64(buffer) { + buffer = int(t.bufferSize) + } + for { + // Re-allocate data object to avoid data collision + buf := make([]byte, buffer) + // Note: ReadFrom receive messages without IP header + n, addr, err := t.conn.ReadFrom(buf) + bufChan <- &RSPacket{buf, addr, err, n} + } + }() for { - // Note: ReadFrom receive messages without IP header - n, addr, err := t.conn.ReadFrom(buf) - - if err != nil { - if strings.HasSuffix(err.Error(), "closed network connection") { + packet := <-bufChan + if packet.err != nil { + if strings.HasSuffix(packet.err.Error(), "closed network connection") { return - } else { - continue } + continue } - if n > 0 { - if t.isValidPacket(buf[:n]) { - t.packetsChan <- t.buildPacket([]byte(addr.(*net.IPAddr).IP), buf[:n], time.Now()) + if packet.n > 0 { + if t.isValidPacket(packet.buf[:packet.n]) { + t.packetsChan <- t.buildPacket([]byte(packet.addr.(*net.IPAddr).IP), packet.buf[:packet.n], time.Now()) } } } @@ -701,16 +712,13 @@ func (t *Listener) isValidPacket(buf []byte) bool { // http://en.wikipedia.org/wiki/Transmission_Control_Protocol destPort := binary.BigEndian.Uint16(buf[2:4]) srcPort := binary.BigEndian.Uint16(buf[0:2]) - // Because RAW_SOCKET can't be bound to port, we have to control it by ourself if destPort == t.port || (t.trackResponse && srcPort == t.port) { // Get the 'data offset' (size of the TCP header in 32-bit words) - dataOffset := (buf[12] & 0xF0) >> 4 - + dataOffset := buf[12] >> 4 // We need only packets with data inside // Check that the buffer is larger than the size of the TCP header if len(buf) > int(dataOffset*4) { - // We should create new buffer because go slices is pointers. So buffer data shoud be immutable. return true } } @@ -880,15 +888,6 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) { } } -func (t *Listener) IsReady() bool { - select { - case <-t.readyCh: - return true - case <-time.After(5 * time.Second): - return false - } -} - // Receiver TCP messages from the listener channel func (t *Listener) Receiver() chan *TCPMessage { return t.messagesChan diff --git a/raw_socket_listener/listener_test.go b/capture/listener_test.go similarity index 94% rename from raw_socket_listener/listener_test.go rename to capture/listener_test.go index 5352e87..578bb58 100644 --- a/raw_socket_listener/listener_test.go +++ b/capture/listener_test.go @@ -1,4 +1,4 @@ -package rawSocket +package capture import ( "bytes" @@ -12,7 +12,7 @@ import ( func TestRawListenerInput(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0, false, false) + listener := NewListener("", "0", EnginePcapFile, true, 10*time.Millisecond, "", "", 0, false, false) defer listener.Close() reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now()) @@ -420,14 +420,6 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket t.Fatal("packetsChan non empty:", listener.packetsChan) } - if len(listener.messagesChan) != 0 { - t.Fatal("messagesChan non empty:", <-listener.messagesChan) - } - - if len(listener.messages) != 0 { - t.Fatal("Messages non empty:", listener.messages) - } - if len(listener.ackAliases) != 0 { t.Fatal("ackAliases non empty:", listener.ackAliases) } @@ -445,24 +437,32 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket } } -func permutation(n int, list []*TCPPacket) []*TCPPacket { - if len(list) == 1 { - return list +// permutation using heap algorithm https://en.wikipedia.org/wiki/Heap%27s_algorithm +func permutation(a []*TCPPacket, f func([]*TCPPacket)) { + n := len(a) + c := make([]int, n) + f(a) + i := 0 + for i < n { + if c[i] < i { + if i&1 != 1 { + a[0], a[i] = a[i], a[0] + } else { + a[c[i]], a[i] = a[i], a[c[i]] + } + f(a) + c[i]++ + i = 0 + } else { + c[i] = 0 + i++ + } } - - k := n % len(list) - - first := []*TCPPacket{list[k]} - next := make([]*TCPPacket, len(list)-1) - - copy(next, append(list[:k], list[k+1:]...)) - - return append(first, permutation(n/len(list), next)...) } // Response comes before Request func TestRawListenerChunkedWrongOrder(t *testing.T) { - listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0, false, false) + listener := NewListener("", "0", EnginePcap, true, 10*time.Second, "", "", 0, false, false) defer listener.Close() reqPacket1 := firstPacket([]byte("POST / HTTP/1.1\r\nTransfer-Encoding: chunked\r\nExpect: 100-continue\r\n\r\n")) @@ -474,12 +474,11 @@ func TestRawListenerChunkedWrongOrder(t *testing.T) { respPacket2 := responsePacket(reqPacket4, []byte("HTTP/1.1 200 OK\r\n\r\n")) - // Should re-construct message from all possible combinations - for i := 0; i < 6*5*4*3*2*1; i++ { - packets := permutation(i, []*TCPPacket{reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket1, respPacket2}) - - testChunkedSequence(t, listener, packets...) + f := func(p []*TCPPacket) { + testChunkedSequence(t, listener, p...) } + // Should re-construct message from all possible combinations + permutation([]*TCPPacket{reqPacket1, reqPacket2, reqPacket3, reqPacket4, respPacket1, respPacket2}, f) } func chunkedPostMessage() []*TCPPacket { diff --git a/raw_socket_listener/tcp_message.go b/capture/tcp_message.go similarity index 99% rename from raw_socket_listener/tcp_message.go rename to capture/tcp_message.go index d09fcca..4ffaf13 100644 --- a/raw_socket_listener/tcp_message.go +++ b/capture/tcp_message.go @@ -1,11 +1,10 @@ -package rawSocket +package capture import ( "bytes" "crypto/sha1" "encoding/binary" "encoding/hex" - "log" "net" "strconv" "strings" @@ -14,8 +13,6 @@ import ( "github.com/buger/goreplay/proto" ) -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 // diff --git a/raw_socket_listener/tcp_message_test.go b/capture/tcp_message_test.go similarity index 99% rename from raw_socket_listener/tcp_message_test.go rename to capture/tcp_message_test.go index e80156d..29170cf 100644 --- a/raw_socket_listener/tcp_message_test.go +++ b/capture/tcp_message_test.go @@ -1,4 +1,4 @@ -package rawSocket +package capture import ( "bytes" diff --git a/raw_socket_listener/tcp_packet.go b/capture/tcp_packet.go similarity index 98% rename from raw_socket_listener/tcp_packet.go rename to capture/tcp_packet.go index 3649a89..6e997a7 100644 --- a/raw_socket_listener/tcp_packet.go +++ b/capture/tcp_packet.go @@ -1,15 +1,12 @@ -package rawSocket +package capture import ( "encoding/binary" - "log" "strconv" "strings" "time" ) -var _ = log.Println - // TCP Flags const ( fFIN = 1 << iota @@ -102,7 +99,6 @@ func (t *TCPPacket) dump() *packet { } copy(packetData[16:], t.Data) - return &packet{ srcIP: packetSrcIP, data: packetData, diff --git a/input_raw.go b/input_raw.go index 47afcef..9fa54d0 100644 --- a/input_raw.go +++ b/input_raw.go @@ -5,8 +5,8 @@ import ( "net" "time" + raw "github.com/buger/goreplay/capture" "github.com/buger/goreplay/proto" - raw "github.com/buger/goreplay/raw_socket_listener" ) // RAWInput used for intercepting traffic for given address @@ -44,9 +44,7 @@ func NewRAWInput(address string, engine int, trackResponse bool, expire time.Dur i.trackResponse = trackResponse i.timestampType = timestampType i.bufferSize = bufferSize - i.listen(address) - i.listener.IsReady() return } @@ -76,7 +74,6 @@ func (i *RAWInput) listen(address string) { Debug("Listening for traffic on: " + address) host, port, err := net.SplitHostPort(address) - if err != nil { log.Fatalf("input-raw: error while parsing address: %s", err) } @@ -90,13 +87,8 @@ func (i *RAWInput) listen(address string) { select { case <-i.quit: return - default: + case i.data <- <-ch: // Receiving TCPMessage object } - - // Receiving TCPMessage object - m := <-ch - - i.data <- m } }() } From 9d71900ad9d528ee6dc4a1ead76b9a6e6d346f31 Mon Sep 17 00:00:00 2001 From: Urban Ishimwe Date: Sat, 6 Jun 2020 10:39:06 +0200 Subject: [PATCH 2/3] listener tests fix --- capture/listener_test.go | 4 ---- 1 file changed, 4 deletions(-) diff --git a/capture/listener_test.go b/capture/listener_test.go index 578bb58..aa47eae 100644 --- a/capture/listener_test.go +++ b/capture/listener_test.go @@ -431,10 +431,6 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket if len(listener.respAliases) != 0 { t.Fatal("respAliases non empty:", listener.respAliases) } - - if len(listener.respWithoutReq) != 0 { - t.Fatal("respWithoutReq non empty:", listener.respWithoutReq) - } } // permutation using heap algorithm https://en.wikipedia.org/wiki/Heap%27s_algorithm From 312d9e8b7dffc1f6884fca3775365be333eac2f7 Mon Sep 17 00:00:00 2001 From: Urban Ishimwe Date: Sat, 6 Jun 2020 10:58:14 +0200 Subject: [PATCH 3/3] optimizing raw socket --- capture/listener.go | 21 ++++++++------------- capture/listener_test.go | 2 +- 2 files changed, 9 insertions(+), 14 deletions(-) diff --git a/capture/listener.go b/capture/listener.go index d3cc70b..b9b851b 100644 --- a/capture/listener.go +++ b/capture/listener.go @@ -74,8 +74,8 @@ type Listener struct { conn net.PacketConn pcapHandles []*pcap.Handle - quit chan bool - readyCh bool + quit chan bool + ready bool } type request struct { @@ -353,7 +353,6 @@ func (t *Listener) readPcap() { } else { inactive.SetSnapLen(65536) } - inactive.SetSnapLen(65536) inactive.SetTimeout(t.messageExpire) inactive.SetPromisc(true) inactive.SetImmediateMode(t.immediateMode) @@ -581,7 +580,7 @@ func (t *Listener) readPcap() { wg.Wait() t.Lock() - t.readyCh = true + t.ready = true t.Unlock() } @@ -597,7 +596,7 @@ func (t *Listener) readPcapFile() { } t.Lock() - t.readyCh = true + t.ready = true t.Unlock() packetSource := gopacket.NewPacketSource(handle, handle.LinkType()) @@ -667,19 +666,15 @@ func (t *Listener) readRAWSocket() { } var bufChan = make(chan *RSPacket, 1000) t.Lock() - t.readyCh = true + t.ready = true t.Unlock() go func() { - buffer := 64 * 1024 - if t.bufferSize > int64(buffer) { - buffer = int(t.bufferSize) - } for { // Re-allocate data object to avoid data collision - buf := make([]byte, buffer) + var buf [64 * 104 * 1024]byte // Note: ReadFrom receive messages without IP header - n, addr, err := t.conn.ReadFrom(buf) - bufChan <- &RSPacket{buf, addr, err, n} + n, addr, err := t.conn.ReadFrom(buf[:]) + bufChan <- &RSPacket{buf[:], addr, err, n} } }() for { diff --git a/capture/listener_test.go b/capture/listener_test.go index aa47eae..7508673 100644 --- a/capture/listener_test.go +++ b/capture/listener_test.go @@ -12,7 +12,7 @@ import ( func TestRawListenerInput(t *testing.T) { var req, resp *TCPMessage - listener := NewListener("", "0", EnginePcapFile, true, 10*time.Millisecond, "", "", 0, false, false) + listener := NewListener("", "0", EnginePcap, true, 10*time.Millisecond, "", "", 0, false, false) defer listener.Close() reqPacket := buildPacket(true, 1, 1, []byte("GET / HTTP/1.1\r\n\r\n"), time.Now())