Update to match upstream conntrack library

This commit is contained in:
Bryan Boreham
2018-08-05 10:42:29 +00:00
parent 5420692a39
commit a29e9fa27a
4 changed files with 81 additions and 95 deletions

View File

@@ -5,7 +5,8 @@ import (
"time"
log "github.com/sirupsen/logrus"
"github.com/weaveworks/scope/probe/endpoint/conntrack"
"github.com/typetypetype/conntrack"
"github.com/weaveworks/scope/probe/endpoint/procspy"
"github.com/weaveworks/scope/probe/process"
"github.com/weaveworks/scope/report"
@@ -54,20 +55,20 @@ func newConnectionTracker(conf connectionTrackerConfig) connectionTracker {
return ct
}
func flowToTuple(f conntrack.Flow) (ft fourTuple) {
func flowToTuple(f conntrack.Conn) (ft fourTuple) {
ft = fourTuple{
f.Original.Layer3.SrcIP.String(),
f.Original.Layer3.DstIP.String(),
uint16(f.Original.Layer4.SrcPort),
uint16(f.Original.Layer4.DstPort),
f.Orig.Src.String(),
f.Orig.Dst.String(),
uint16(f.Orig.SrcPort),
uint16(f.Orig.DstPort),
}
// Handle DNAT-ed connections in the initial state
if !f.Original.Layer3.DstIP.Equal(f.Reply.Layer3.SrcIP) {
if !f.Orig.Dst.Equal(f.Reply.Src) {
ft = fourTuple{
f.Reply.Layer3.DstIP.String(),
f.Reply.Layer3.SrcIP.String(),
uint16(f.Reply.Layer4.DstPort),
uint16(f.Reply.Layer4.SrcPort),
f.Reply.Dst.String(),
f.Reply.Src.String(),
uint16(f.Reply.DstPort),
uint16(f.Reply.SrcPort),
}
}
return ft
@@ -119,7 +120,7 @@ func (t *connectionTracker) ReportConnections(rpt *report.Report) {
// consult the flowWalker for short-lived (conntracked) connections
seenTuples := map[string]fourTuple{}
t.flowWalker.walkFlows(func(f conntrack.Flow, alive bool) {
t.flowWalker.walkFlows(func(f conntrack.Conn, alive bool) {
tuple := flowToTuple(f)
seenTuples[tuple.key()] = tuple
t.addConnection(rpt, false, tuple, "", nil, nil)
@@ -136,7 +137,7 @@ func (t *connectionTracker) existingFlows() map[string]fourTuple {
// log.Warnf("Not using conntrack: disabled")
} else if err := IsConntrackSupported(t.conf.ProcRoot); err != nil {
log.Warnf("Not using conntrack: not supported by the kernel: %s", err)
} else if existingFlows, err := conntrack.Established(t.conf.BufferSize); err != nil {
} else if existingFlows, err := conntrack.ConnectionsSize(t.conf.BufferSize); err != nil {
log.Errorf("conntrack existingConnections error: %v", err)
} else {
for _, f := range existingFlows {

View File

@@ -10,8 +10,7 @@ import (
"time"
log "github.com/sirupsen/logrus"
"github.com/weaveworks/scope/probe/endpoint/conntrack"
"github.com/typetypetype/conntrack"
)
const (
@@ -22,21 +21,21 @@ const (
// flowWalker is something that maintains flows, and provides an accessor
// method to walk them.
type flowWalker interface {
walkFlows(f func(conntrack.Flow, bool))
walkFlows(f func(conntrack.Conn, bool))
stop()
}
type nilFlowWalker struct{}
func (n nilFlowWalker) stop() {}
func (n nilFlowWalker) walkFlows(f func(conntrack.Flow, bool)) {}
func (n nilFlowWalker) walkFlows(f func(conntrack.Conn, bool)) {}
// conntrackWalker uses the conntrack command to track network connections and
// implement flowWalker.
type conntrackWalker struct {
sync.Mutex
activeFlows map[uint32]conntrack.Flow // active flows in state != TIME_WAIT
bufferedFlows []conntrack.Flow // flows coming out of activeFlows spend 1 walk cycle here
activeFlows map[uint32]conntrack.Conn // active flows in state != TIME_WAIT
bufferedFlows []conntrack.Conn // flows coming out of activeFlows spend 1 walk cycle here
bufferSize int
natOnly bool
quit chan struct{}
@@ -51,7 +50,7 @@ func newConntrackFlowWalker(useConntrack bool, procRoot string, bufferSize int,
return nilFlowWalker{}
}
result := &conntrackWalker{
activeFlows: map[uint32]conntrack.Flow{},
activeFlows: map[uint32]conntrack.Conn{},
bufferSize: bufferSize,
natOnly: natOnly,
quit: make(chan struct{}),
@@ -100,7 +99,7 @@ func (c *conntrackWalker) clearFlows() {
c.bufferedFlows = append(c.bufferedFlows, f)
}
c.activeFlows = map[uint32]conntrack.Flow{}
c.activeFlows = map[uint32]conntrack.Conn{}
}
func logPipe(prefix string, reader io.Reader) {
@@ -123,9 +122,9 @@ func (c *conntrackWalker) run() {
c.handleFlow(flow, true)
}
events, stop, err := conntrack.Follow(c.bufferSize)
events, stop, err := conntrack.FollowSize(c.bufferSize, conntrack.NF_NETLINK_CONNTRACK_UPDATE|conntrack.NF_NETLINK_CONNTRACK_DESTROY)
if err != nil {
log.Errorf("conntract Follow error: %v", err)
log.Errorf("conntrack Follow error: %v", err)
return
}
@@ -157,10 +156,10 @@ func (c *conntrackWalker) run() {
}
}
func (c *conntrackWalker) existingConnections() ([]conntrack.Flow, error) {
flows, err := conntrack.Established(c.bufferSize)
func (c *conntrackWalker) existingConnections() ([]conntrack.Conn, error) {
flows, err := conntrack.ConnectionsSize(c.bufferSize)
if err != nil {
return []conntrack.Flow{}, err
return []conntrack.Conn{}, err
}
return flows, nil
}
@@ -171,7 +170,7 @@ func (c *conntrackWalker) stop() {
close(c.quit)
}
func (c *conntrackWalker) handleFlow(f conntrack.Flow, forceAdd bool) {
func (c *conntrackWalker) handleFlow(f conntrack.Conn, forceAdd bool) {
c.Lock()
defer c.Unlock()
@@ -183,15 +182,15 @@ func (c *conntrackWalker) handleFlow(f conntrack.Flow, forceAdd bool) {
// incomplete or wrong. See #1462.
switch {
case forceAdd || f.MsgType == conntrack.NfctMsgUpdate:
if f.State != conntrack.TCPStateTimeWait {
c.activeFlows[f.ID] = f
} else if _, ok := c.activeFlows[f.ID]; ok {
delete(c.activeFlows, f.ID)
if f.TCPState != "TIME_WAIT" {
c.activeFlows[f.CtId] = f
} else if _, ok := c.activeFlows[f.CtId]; ok {
delete(c.activeFlows, f.CtId)
c.bufferedFlows = append(c.bufferedFlows, f)
}
case f.MsgType == conntrack.NfctMsgDestroy:
if active, ok := c.activeFlows[f.ID]; ok {
delete(c.activeFlows, f.ID)
if active, ok := c.activeFlows[f.CtId]; ok {
delete(c.activeFlows, f.CtId)
c.bufferedFlows = append(c.bufferedFlows, active)
}
}
@@ -199,7 +198,7 @@ func (c *conntrackWalker) handleFlow(f conntrack.Flow, forceAdd bool) {
// walkFlows calls f with all active flows and flows that have come and gone
// since the last call to walkFlows
func (c *conntrackWalker) walkFlows(f func(conntrack.Flow, bool)) {
func (c *conntrackWalker) walkFlows(f func(conntrack.Conn, bool)) {
c.Lock()
defer c.Unlock()
for _, flow := range c.activeFlows {

View File

@@ -4,7 +4,8 @@ import (
"net"
"strconv"
"github.com/weaveworks/scope/probe/endpoint/conntrack"
"github.com/typetypetype/conntrack"
"github.com/weaveworks/scope/report"
)
@@ -27,21 +28,21 @@ func makeNATMapper(fw flowWalker) natMapper {
return natMapper{fw}
}
func toMapping(f conntrack.Flow) *endpointMapping {
func toMapping(f conntrack.Conn) *endpointMapping {
var mapping endpointMapping
if f.Original.Layer3.SrcIP.Equal(f.Reply.Layer3.DstIP) {
if f.Orig.Src.Equal(f.Reply.Dst) {
mapping = endpointMapping{
originalIP: f.Reply.Layer3.SrcIP,
originalPort: f.Reply.Layer4.SrcPort,
rewrittenIP: f.Original.Layer3.DstIP,
rewrittenPort: f.Original.Layer4.DstPort,
originalIP: f.Reply.Src,
originalPort: f.Reply.SrcPort,
rewrittenIP: f.Orig.Dst,
rewrittenPort: f.Orig.DstPort,
}
} else {
mapping = endpointMapping{
originalIP: f.Original.Layer3.SrcIP,
originalPort: f.Original.Layer4.SrcPort,
rewrittenIP: f.Reply.Layer3.DstIP,
rewrittenPort: f.Reply.Layer4.DstPort,
originalIP: f.Orig.Src,
originalPort: f.Orig.SrcPort,
rewrittenIP: f.Reply.Dst,
rewrittenPort: f.Reply.DstPort,
}
}
@@ -51,7 +52,7 @@ func toMapping(f conntrack.Flow) *endpointMapping {
// applyNAT duplicates Nodes in the endpoint topology of a report, based on
// the NAT table.
func (n natMapper) applyNAT(rpt report.Report, scope string) {
n.flowWalker.walkFlows(func(f conntrack.Flow, _ bool) {
n.flowWalker.walkFlows(func(f conntrack.Conn, _ bool) {
mapping := toMapping(f)
realEndpointPort := strconv.Itoa(int(mapping.originalPort))

View File

@@ -5,18 +5,19 @@ import (
"syscall"
"testing"
"github.com/typetypetype/conntrack"
"github.com/weaveworks/common/mtime"
"github.com/weaveworks/common/test"
"github.com/weaveworks/scope/probe/endpoint/conntrack"
"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/test/reflect"
)
type mockFlowWalker struct {
flows []conntrack.Flow
flows []conntrack.Conn
}
func (m *mockFlowWalker) walkFlows(f func(f conntrack.Flow, active bool)) {
func (m *mockFlowWalker) walkFlows(f func(f conntrack.Conn, active bool)) {
for _, flow := range m.flows {
f(flow, true)
}
@@ -42,35 +43,27 @@ func TestNat(t *testing.T) {
// from the PoV of host1
{
f := conntrack.Flow{
f := conntrack.Conn{
MsgType: conntrack.NfctMsgUpdate,
Original: conntrack.Meta{
Layer3: conntrack.Layer3{
SrcIP: host2,
DstIP: host1,
},
Layer4: conntrack.Layer4{
SrcPort: 22222,
DstPort: 80,
Proto: syscall.IPPROTO_TCP,
},
Orig: conntrack.Tuple{
Src: host2,
Dst: host1,
SrcPort: 22222,
DstPort: 80,
Proto: syscall.IPPROTO_TCP,
},
Reply: conntrack.Meta{
Layer3: conntrack.Layer3{
SrcIP: c1,
DstIP: host2,
},
Layer4: conntrack.Layer4{
SrcPort: 80,
DstPort: 22222,
Proto: syscall.IPPROTO_TCP,
},
Reply: conntrack.Tuple{
Src: c1,
Dst: host2,
SrcPort: 80,
DstPort: 22222,
Proto: syscall.IPPROTO_TCP,
},
ID: 1,
CtId: 1,
}
ct := &mockFlowWalker{
flows: []conntrack.Flow{f},
flows: []conntrack.Conn{f},
}
have := report.MakeReport()
@@ -94,34 +87,26 @@ func TestNat(t *testing.T) {
// form the PoV of host2
{
f := conntrack.Flow{
f := conntrack.Conn{
MsgType: conntrack.NfctMsgUpdate,
Original: conntrack.Meta{
Layer3: conntrack.Layer3{
SrcIP: c2,
DstIP: host1,
},
Layer4: conntrack.Layer4{
SrcPort: 22222,
DstPort: 80,
Proto: syscall.IPPROTO_TCP,
},
Orig: conntrack.Tuple{
Src: c2,
Dst: host1,
SrcPort: 22222,
DstPort: 80,
Proto: syscall.IPPROTO_TCP,
},
Reply: conntrack.Meta{
Layer3: conntrack.Layer3{
SrcIP: host1,
DstIP: host2,
},
Layer4: conntrack.Layer4{
SrcPort: 80,
DstPort: 22223,
Proto: syscall.IPPROTO_TCP,
},
Reply: conntrack.Tuple{
Src: host1,
Dst: host2,
SrcPort: 80,
DstPort: 22223,
Proto: syscall.IPPROTO_TCP,
},
ID: 2,
CtId: 2,
}
ct := &mockFlowWalker{
flows: []conntrack.Flow{f},
flows: []conntrack.Conn{f},
}
have := report.MakeReport()