Revert "Make packet proccessing multi threaded"

This reverts commit 11d61dcb4d.

This change was causing some packet loss
This commit is contained in:
Leonid Bugaev
2021-08-26 19:10:59 +03:00
parent 05ed821297
commit 1025b8bdb2
+23 -57
View File
@@ -7,7 +7,6 @@ import (
"net"
"reflect"
"sort"
"sync"
"time"
"unsafe"
)
@@ -65,7 +64,6 @@ type Message struct {
packets []*Packet
parser *MessageParser
feedback interface{}
Idx uint16
Stats
}
@@ -200,8 +198,7 @@ type HintStart func(*Packet) (IsRequest, IsOutgoing bool)
// MessageParser holds data of all tcp messages in progress(still receiving/sending packets).
// message is identified by its source port and dst port, and last 4bytes of src IP.
type MessageParser struct {
m []map[uint64]*Message
mL []sync.RWMutex
m map[uint64]*Message
messageExpire time.Duration // the maximum time to wait for the final packet, minimum is 100ms
allowIncompete bool
@@ -229,24 +226,18 @@ func NewMessageParser(messages chan *Message, ports []uint16, ips []net.IP, mess
parser.packets = make(chan *PcapPacket, 10000)
if messages == nil {
messages = make(chan *Message, 100)
messages = make(chan *Message, 1000)
}
parser.messages = messages
parser.m = make(map[uint64]*Message)
parser.ticker = time.NewTicker(time.Millisecond * 100)
parser.close = make(chan struct{}, 1)
parser.ports = ports
parser.ips = ips
for i := 0; i < 10; i++ {
parser.m = append(parser.m, make(map[uint64]*Message))
parser.mL = append(parser.mL, sync.RWMutex{})
}
for i := 0; i < 10; i++ {
go parser.wait(i)
}
go parser.wait()
return parser
}
@@ -258,7 +249,7 @@ func (parser *MessageParser) PacketHandler(packet *PcapPacket) {
parser.packets <- packet
}
func (parser *MessageParser) wait(index int) {
func (parser *MessageParser) wait() {
var (
now time.Time
)
@@ -267,7 +258,7 @@ func (parser *MessageParser) wait(index int) {
case pckt := <-parser.packets:
parser.processPacket(parser.parsePacket(pckt))
case now = <-parser.ticker.C:
parser.timer(now, index)
parser.timer(now)
case <-parser.close:
parser.ticker.Stop()
// parser.Close should wait for this function to return
@@ -307,28 +298,19 @@ func (parser *MessageParser) processPacket(pckt *Packet) {
// Trying to build unique hash, but there is small chance of collision
// No matter if it is request or response, all packets in the same message have same
mID := pckt.MessageID()
mIDX := pckt.SrcPort % 10
parser.mL[mIDX].Lock()
m, ok := parser.m[mIDX][mID]
if !ok {
parser.mL[mIDX].Unlock()
mIDX = pckt.DstPort % 10
parser.mL[mIDX].Lock()
m, ok = parser.m[mIDX][mID]
if !ok {
parser.mL[mIDX].Unlock()
}
}
m, ok := parser.m[pckt.MessageID()]
switch {
case ok:
if m.Direction == DirUnknown {
if in, out := parser.Start(pckt); in || out {
if in {
m.Direction = DirIncoming
} else {
m.Direction = DirOutcoming
}
}
}
parser.addPacket(m, pckt)
parser.mL[mIDX].Unlock()
return
case pckt.Direction == DirUnknown && parser.Start != nil:
if in, out := parser.Start(pckt); in || out {
@@ -340,27 +322,14 @@ func (parser *MessageParser) processPacket(pckt *Packet) {
}
}
if pckt.Direction == DirIncoming {
mIDX = pckt.SrcPort % 10
} else {
mIDX = pckt.DstPort % 10
}
parser.mL[mIDX].Lock()
m = new(Message)
m.Direction = pckt.Direction
m.SrcAddr = pckt.SrcIP.String()
m.DstAddr = pckt.DstIP.String()
parser.m[mIDX][mID] = m
m.Idx = mIDX
parser.m[pckt.MessageID()] = m
m.Start = pckt.Timestamp
m.parser = parser
parser.addPacket(m, pckt)
parser.mL[mIDX].Unlock()
}
func (parser *MessageParser) addPacket(m *Message, pckt *Packet) bool {
@@ -387,7 +356,7 @@ func (parser *MessageParser) Read() *Message {
func (parser *MessageParser) Emit(m *Message) {
stats.Add("message_count", 1)
delete(parser.m[m.Idx], m.packets[0].MessageID())
delete(parser.m, m.packets[0].MessageID())
parser.messages <- m
}
@@ -398,14 +367,13 @@ func GetUnexportedField(field reflect.Value) interface{} {
var failMsg int
func (parser *MessageParser) timer(now time.Time, index int) {
func (parser *MessageParser) timer(now time.Time) {
packetLen = 0
parser.mL[index].Lock()
packetQueueLen.Set(int64(len(parser.packets)))
messageQueueLen.Set(int64(len(parser.m[index])))
messageQueueLen.Set(int64(len(parser.m)))
for _, m := range parser.m[index] {
for _, m := range parser.m {
if now.Sub(m.End) > parser.messageExpire {
m.TimedOut = true
stats.Add("message_timeout_count", 1)
@@ -414,11 +382,9 @@ func (parser *MessageParser) timer(now time.Time, index int) {
parser.Emit(m)
}
delete(parser.m[index], m.packets[0].MessageID())
delete(parser.m, m.packets[0].MessageID())
}
}
parser.mL[index].Unlock()
}
func (parser *MessageParser) Close() error {