Apply FMT

This commit is contained in:
Leonid Bugaev
2016-05-11 16:53:23 +05:00
parent 47bb312b73
commit 768c5f9726
7 changed files with 56 additions and 59 deletions
+1 -1
View File
@@ -81,7 +81,7 @@ func (i *RAWInput) listen(address string) {
} }
// Receiving TCPMessage object // Receiving TCPMessage object
m := <- ch m := <-ch
i.data <- m i.data <- m
} }
+8 -10
View File
@@ -31,16 +31,15 @@ func TestRAWInput(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
origin := &http.Server{ origin := &http.Server{
Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}), Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}),
ReadTimeout: 10 * time.Second, ReadTimeout: 10 * time.Second,
WriteTimeout: 10 * time.Second, WriteTimeout: 10 * time.Second,
} }
go origin.Serve(listener) go origin.Serve(listener)
defer listener.Close() defer listener.Close()
originAddr := listener.Addr().String() originAddr := listener.Addr().String()
var respCounter, reqCounter int64 var respCounter, reqCounter int64
input := NewRAWInput(originAddr, EnginePcap, testRawExpire) input := NewRAWInput(originAddr, EnginePcap, testRawExpire)
@@ -63,7 +62,7 @@ func TestRAWInput(t *testing.T) {
Plugins.Inputs = []io.Reader{input} Plugins.Inputs = []io.Reader{input}
Plugins.Outputs = []io.Writer{output} Plugins.Outputs = []io.Writer{output}
client := NewHTTPClient("http://" + listener.Addr().String(), &HTTPClientConfig{}) client := NewHTTPClient("http://"+listener.Addr().String(), &HTTPClientConfig{})
go Start(quit) go Start(quit)
@@ -87,16 +86,15 @@ func TestRAWInputIPv6(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
origin := &http.Server{ origin := &http.Server{
Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}), Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}),
ReadTimeout: 10 * time.Second, ReadTimeout: 10 * time.Second,
WriteTimeout: 10 * time.Second, WriteTimeout: 10 * time.Second,
} }
go origin.Serve(listener) go origin.Serve(listener)
defer listener.Close() defer listener.Close()
originAddr := listener.Addr().String() originAddr := listener.Addr().String()
var respCounter, reqCounter int64 var respCounter, reqCounter int64
input := NewRAWInput(originAddr, EnginePcap, testRawExpire) input := NewRAWInput(originAddr, EnginePcap, testRawExpire)
@@ -119,7 +117,7 @@ func TestRAWInputIPv6(t *testing.T) {
Plugins.Inputs = []io.Reader{input} Plugins.Inputs = []io.Reader{input}
Plugins.Outputs = []io.Writer{output} Plugins.Outputs = []io.Writer{output}
client := NewHTTPClient("http://" + listener.Addr().String(), &HTTPClientConfig{}) client := NewHTTPClient("http://"+listener.Addr().String(), &HTTPClientConfig{})
go Start(quit) go Start(quit)
+14 -14
View File
@@ -61,8 +61,8 @@ type Listener struct {
messageExpire time.Duration messageExpire time.Duration
conn net.PacketConn conn net.PacketConn
quit chan bool quit chan bool
readyCh chan bool readyCh chan bool
} }
@@ -130,10 +130,10 @@ func (t *Listener) listen() {
t.conn.Close() t.conn.Close()
} }
return return
case data := <- t.packetsChan: case data := <-t.packetsChan:
packet := ParseTCPPacket(data[:16], data[16:]) packet := ParseTCPPacket(data[:16], data[16:])
t.processTCPPacket(packet) t.processTCPPacket(packet)
case <- gcTicker: case <-gcTicker:
now := time.Now() now := time.Now()
// Dispatch requests before responses // Dispatch requests before responses
@@ -175,13 +175,13 @@ func (t *Listener) dispatchMessage(message *TCPMessage) {
if respID, ok := t.respWithoutReq[message.ResponseAck]; ok { if respID, ok := t.respWithoutReq[message.ResponseAck]; ok {
if resp, rok := t.messages[respID]; rok { if resp, rok := t.messages[respID]; rok {
// if resp.AssocMessage == nil { // if resp.AssocMessage == nil {
// log.Println("FOUND RESPONSE") // log.Println("FOUND RESPONSE")
resp.AssocMessage = message resp.AssocMessage = message
message.AssocMessage = resp message.AssocMessage = resp
if resp.IsFinished() { if resp.IsFinished() {
defer t.dispatchMessage(resp) defer t.dispatchMessage(resp)
} }
// } // }
} }
} }
@@ -333,7 +333,7 @@ func (t *Listener) readPcap() {
// We need only packets with data inside // We need only packets with data inside
// Check that the buffer is larger than the size of the TCP header // Check that the buffer is larger than the size of the TCP header
if len(data) > int(dataOffset*4) { if len(data) > int(dataOffset*4) {
newBuf := make([]byte, len(data) + 16) newBuf := make([]byte, len(data)+16)
copy(newBuf[:16], srcIP) copy(newBuf[:16], srcIP)
copy(newBuf[16:], data) copy(newBuf[16:], data)
@@ -375,7 +375,7 @@ func (t *Listener) readRAWSocket() {
if n > 0 { if n > 0 {
if t.isValidPacket(buf[:n]) { if t.isValidPacket(buf[:n]) {
newBuf := make([]byte, n + 4) newBuf := make([]byte, n+4)
copy(newBuf[16:], buf[:n]) copy(newBuf[16:], buf[:n])
copy(newBuf[:16], []byte(addr.(*net.IPAddr).IP)) copy(newBuf[:16], []byte(addr.(*net.IPAddr).IP))
@@ -550,9 +550,9 @@ func (t *Listener) processTCPPacket(packet *TCPPacket) {
func (t *Listener) IsReady() bool { func (t *Listener) IsReady() bool {
select { select {
case <- t.readyCh: case <-t.readyCh:
return true return true
case <- time.After(5 * time.Second): case <-time.After(5 * time.Second):
return false return false
} }
} }
+18 -19
View File
@@ -3,10 +3,10 @@ package rawSocket
import ( import (
"bytes" "bytes"
"log" "log"
"testing"
"time"
"math/rand" "math/rand"
"sync/atomic" "sync/atomic"
"testing"
"time"
) )
func TestRawListenerInput(t *testing.T) { func TestRawListenerInput(t *testing.T) {
@@ -262,7 +262,7 @@ func testChunkedSequence(t *testing.T, listener *Listener, packets ...*TCPPacket
} }
if len(listener.messagesChan) != 0 { if len(listener.messagesChan) != 0 {
t.Fatal("messagesChan non empty:", <- listener.messagesChan) t.Fatal("messagesChan non empty:", <-listener.messagesChan)
} }
if len(listener.messages) != 0 { if len(listener.messages) != 0 {
@@ -331,14 +331,13 @@ func TestRawListenerChunkedWrongOrder(t *testing.T) {
} }
} }
func chunkedPostMessage() []*TCPPacket { func chunkedPostMessage() []*TCPPacket {
ack := uint32(rand.Int63()) ack := uint32(rand.Int63())
seq := 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"))
// Packet with data have different Seq // Packet with data have different Seq
reqPacket2 := buildPacket(true, ack, seq + 47, []byte("1\r\na\r\n")) 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")) 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")) reqPacket4 := buildPacket(true, ack, reqPacket3.Seq+5, []byte("0\r\n\r\n"))
@@ -359,13 +358,13 @@ func postMessage() []*TCPPacket {
rand.Read(data) rand.Read(data)
head := []byte("POST / HTTP/1.1\r\nContent-Length: 9958\r\n\r\n") head := []byte("POST / HTTP/1.1\r\nContent-Length: 9958\r\n\r\n")
for i, _ := range head { for i := range head {
data[i] = head[i] data[i] = head[i]
} }
return []*TCPPacket{ return []*TCPPacket{
buildPacket(true, ack, seq, data), buildPacket(true, ack, seq, data),
buildPacket(false, seq + uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n")), buildPacket(false, seq+uint32(len(data)), seq2, []byte("HTTP/1.1 200 OK\r\n")),
} }
} }
@@ -376,7 +375,7 @@ func getMessage() []*TCPPacket {
return []*TCPPacket{ return []*TCPPacket{
buildPacket(true, ack, seq, []byte("GET / HTTP/1.1\r\n\r\n")), 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")), buildPacket(false, seq+18, seq2, []byte("HTTP/1.1 200 OK\r\n")),
} }
} }
@@ -387,22 +386,22 @@ func TestRawListenerBench(t *testing.T) {
// Should re-construct message from all possible combinations // Should re-construct message from all possible combinations
for i := 0; i < 1000; i++ { for i := 0; i < 1000; i++ {
go func(){ go func() {
for j := 0; j < 100; j++ { for j := 0; j < 100; j++ {
var packets []*TCPPacket var packets []*TCPPacket
if j % 5 == 0 { if j%5 == 0 {
packets = chunkedPostMessage() packets = chunkedPostMessage()
} else if j % 3 == 0 { } else if j%3 == 0 {
packets = postMessage() packets = postMessage()
} else { } else {
packets = getMessage() packets = getMessage()
} }
for _, p := range packets { for _, p := range packets {
// Randomly drop packets // Randomly drop packets
if (i + j) % 5 == 0 { if (i+j)%5 == 0 {
if rand.Int63() % 3 == 0 { if rand.Int63()%3 == 0 {
continue continue
} }
} }
@@ -422,11 +421,11 @@ func TestRawListenerBench(t *testing.T) {
for { for {
select { select {
case <- ch: case <-ch:
atomic.AddInt32(&count, 1) atomic.AddInt32(&count, 1)
case <-time.After(2000 * time.Millisecond): case <-time.After(2000 * time.Millisecond):
log.Println("Emitted 200000 messages, captured: ", count, len(l.ackAliases), len(l.seqWithData), len(l.respAliases), len(l.respWithoutReq), len(l.packetsChan)) log.Println("Emitted 200000 messages, captured: ", count, len(l.ackAliases), len(l.seqWithData), len(l.respAliases), len(l.respWithoutReq), len(l.packetsChan))
return return
} }
} }
} }
+9 -9
View File
@@ -3,8 +3,8 @@ package rawSocket
import ( import (
"bytes" "bytes"
"crypto/sha1" "crypto/sha1"
"encoding/hex"
"encoding/binary" "encoding/binary"
"encoding/hex"
"github.com/buger/gor/proto" "github.com/buger/gor/proto"
"log" "log"
"strconv" "strconv"
@@ -20,12 +20,12 @@ var _ = log.Println
// Message can be compiled from unique packets with same message_id which sorted by sequence // Message can be compiled from unique packets with same message_id which sorted by sequence
// Message is received if we didn't receive any packets for 2000ms // Message is received if we didn't receive any packets for 2000ms
type TCPMessage struct { type TCPMessage struct {
Seq uint32 Seq uint32
Ack uint32 Ack uint32
ResponseAck uint32 ResponseAck uint32
ResponseID tcpID ResponseID tcpID
DataAck uint32 DataAck uint32
DataSeq uint32 DataSeq uint32
AssocMessage *TCPMessage AssocMessage *TCPMessage
Start time.Time Start time.Time
@@ -210,11 +210,11 @@ func (t *TCPMessage) UpdateResponseAck() uint32 {
copy(t.ResponseID[4:], lastPacket.Raw[2:4]) // Src port copy(t.ResponseID[4:], lastPacket.Raw[2:4]) // Src port
copy(t.ResponseID[6:], lastPacket.Raw[0:2]) // Dest port copy(t.ResponseID[6:], lastPacket.Raw[0:2]) // Dest port
binary.BigEndian.PutUint32(t.ResponseID[8:12], t.ResponseAck) binary.BigEndian.PutUint32(t.ResponseID[8:12], t.ResponseAck)
} }
return t.ResponseAck return t.ResponseAck
} }
func (t *TCPMessage) ID() tcpID { func (t *TCPMessage) ID() tcpID {
return t.packets[0].ID return t.packets[0].ID
} }
+1 -1
View File
@@ -2,9 +2,9 @@ package rawSocket
import ( import (
"bytes" "bytes"
"encoding/binary"
_ "log" _ "log"
"testing" "testing"
"encoding/binary"
) )
func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPacket) { func buildPacket(isIncoming bool, Ack, Seq uint32, Data []byte) (packet *TCPPacket) {
+5 -5
View File
@@ -31,10 +31,10 @@ type TCPPacket struct {
OrigAck uint32 OrigAck uint32
DataOffset uint8 DataOffset uint8
Raw []byte Raw []byte
Data []byte Data []byte
Addr []byte Addr []byte
ID tcpID ID tcpID
} }
// ParseTCPPacket takes address and tcp payload and returns parsed TCPPacket // ParseTCPPacket takes address and tcp payload and returns parsed TCPPacket
@@ -49,8 +49,8 @@ func ParseTCPPacket(addr []byte, data []byte) (p *TCPPacket) {
func (p *TCPPacket) GenID() { func (p *TCPPacket) GenID() {
copy(p.ID[:16], p.Addr) copy(p.ID[:16], p.Addr)
copy(p.ID[16:], p.Raw[0:2]) // Src port copy(p.ID[16:], p.Raw[0:2]) // Src port
copy(p.ID[18:], p.Raw[2:4]) // Dest port copy(p.ID[18:], p.Raw[2:4]) // Dest port
copy(p.ID[20:], p.Raw[8:12]) // Ack copy(p.ID[20:], p.Raw[8:12]) // Ack
} }
@@ -73,7 +73,7 @@ func (t *TCPPacket) ParseBasic() {
} }
func (t *TCPPacket) Dump() []byte { func (t *TCPPacket) Dump() []byte {
buf := make([]byte, len(t.Data) + 16 + 16) buf := make([]byte, len(t.Data)+16+16)
tcpBuf := buf[16:] tcpBuf := buf[16:]
binary.BigEndian.PutUint16(tcpBuf[2:4], t.DestPort) binary.BigEndian.PutUint16(tcpBuf[2:4], t.DestPort)