From a29e9fa27ac3ae4b52b3f506708e7901242f12ac Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 5 Aug 2018 10:42:29 +0000 Subject: [PATCH] Update to match upstream conntrack library --- probe/endpoint/connection_tracker.go | 27 ++++----- probe/endpoint/conntrack.go | 41 +++++++------- probe/endpoint/nat.go | 25 +++++---- probe/endpoint/nat_internal_test.go | 83 ++++++++++++---------------- 4 files changed, 81 insertions(+), 95 deletions(-) diff --git a/probe/endpoint/connection_tracker.go b/probe/endpoint/connection_tracker.go index ee2e749a5..73c0861df 100644 --- a/probe/endpoint/connection_tracker.go +++ b/probe/endpoint/connection_tracker.go @@ -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 { diff --git a/probe/endpoint/conntrack.go b/probe/endpoint/conntrack.go index de9607e25..02a1dc93d 100644 --- a/probe/endpoint/conntrack.go +++ b/probe/endpoint/conntrack.go @@ -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 { diff --git a/probe/endpoint/nat.go b/probe/endpoint/nat.go index aa87ab838..d5eb71001 100644 --- a/probe/endpoint/nat.go +++ b/probe/endpoint/nat.go @@ -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)) diff --git a/probe/endpoint/nat_internal_test.go b/probe/endpoint/nat_internal_test.go index 57a09cb9c..27a1cc624 100644 --- a/probe/endpoint/nat_internal_test.go +++ b/probe/endpoint/nat_internal_test.go @@ -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()