From 0aadf6447b1cce85bacf85efb14902103706d726 Mon Sep 17 00:00:00 2001 From: Peter Bourgon Date: Fri, 31 Jul 2015 19:47:55 +0200 Subject: [PATCH] Revert to correct edge construction Another implicit invariant in the data model is that edges are always of the form (local -> remote). That is, the source of an edge must always be a node that originates from within Scope's domain of visibility. This was evident by the presence of ingress and egress fields in edge/aggregate metadata. When building the sniffer, I accidentally and incorrectly violated this invariant, by constructing distinct edges for (local -> remote) and (remote -> local), and collapsing ingress and egress byte counts to a single scalar. I experienced a variety of subtle undefined behavior as a result. See #339. This change reverts to the old, correct methodology. Consequently the sniffer needs to be able to find out which side of the sniffed packet is local v. remote, and to do that it needs access to local networks. I moved the discovery from the probe/host package into probe/main.go. As part of that work I discovered that package report also maintains its own, independent "cache" of local networks. Except it contains only the (optional) Docker bridge network, if it's been populated by the probe, and it's only used by the report.Make{Endpoint,Address}NodeID constructors to scope local addresses. Normally, scoping happens during rendering, and only for pseudo nodes -- see current LeafMap Render localNetworks. This is pretty convoluted and should be either be made consistent or heavily commented. --- app/api_topology_test.go | 4 +- probe/host/reporter.go | 33 +++---- probe/host/reporter_test.go | 32 ++++--- probe/main.go | 16 +++- probe/sniff/sniffer.go | 128 +++++++++++++++++++-------- probe/sniff/sniffer_internal_test.go | 10 ++- probe/sniff/sniffer_test.go | 18 ++-- render/detailed_node.go | 7 +- render/detailed_node_test.go | 2 +- render/expected/expected.go | 36 ++++---- render/render_test.go | 8 +- report/merge.go | 6 +- report/merge_test.go | 38 ++++---- report/networks.go | 7 +- report/topology.go | 7 +- test/report_fixture.go | 28 +++--- 16 files changed, 225 insertions(+), 155 deletions(-) diff --git a/app/api_topology_test.go b/app/api_topology_test.go index 590b9fa02..61fa8fa84 100644 --- a/app/api_topology_test.go +++ b/app/api_topology_test.go @@ -96,8 +96,8 @@ func TestAPITopologyApplications(t *testing.T) { t.Fatalf("JSON parse error: %s", err) } if want, have := (report.EdgeMetadata{ - PacketCount: newu64(100), - ByteCount: newu64(10), + PacketCount: newu64(10), + EgressByteCount: newu64(100), }), edge.Metadata; !reflect.DeepEqual(want, have) { t.Error(test.Diff(want, have)) } diff --git a/probe/host/reporter.go b/probe/host/reporter.go index 6694e96f7..a242e6a6f 100644 --- a/probe/host/reporter.go +++ b/probe/host/reporter.go @@ -1,7 +1,6 @@ package host import ( - "net" "runtime" "strings" "time" @@ -9,7 +8,7 @@ import ( "github.com/weaveworks/scope/report" ) -// Keys for use in NodeMetadata +// Keys for use in NodeMetadata. const ( Timestamp = "ts" HostName = "host_name" @@ -20,30 +19,31 @@ const ( Uptime = "uptime" ) -// Exposed for testing +// Exposed for testing. const ( ProcUptime = "/proc/uptime" ProcLoad = "/proc/loadavg" ) -// Exposed for testing +// Exposed for testing. var ( - InterfaceAddrs = net.InterfaceAddrs - Now = func() string { return time.Now().UTC().Format(time.RFC3339Nano) } + Now = func() string { return time.Now().UTC().Format(time.RFC3339Nano) } ) // Reporter generates Reports containing the host topology. type Reporter struct { - hostID string - hostName string + hostID string + hostName string + localNets report.Networks } // NewReporter returns a Reporter which produces a report containing host // topology for this host. -func NewReporter(hostID, hostName string) *Reporter { +func NewReporter(hostID, hostName string, localNets report.Networks) *Reporter { return &Reporter{ - hostID: hostID, - hostName: hostName, + hostID: hostID, + hostName: hostName, + localNets: localNets, } } @@ -54,15 +54,8 @@ func (r *Reporter) Report() (report.Report, error) { localCIDRs []string ) - localNets, err := InterfaceAddrs() - if err != nil { - return rep, err - } - for _, localNet := range localNets { - // Not all networks are IP networks. - if ipNet, ok := localNet.(*net.IPNet); ok { - localCIDRs = append(localCIDRs, ipNet.String()) - } + for _, localNet := range r.localNets { + localCIDRs = append(localCIDRs, localNet.String()) } uptime, err := GetUptime() diff --git a/probe/host/reporter_test.go b/probe/host/reporter_test.go index f4127a577..a3fa0df09 100644 --- a/probe/host/reporter_test.go +++ b/probe/host/reporter_test.go @@ -12,38 +12,37 @@ import ( "github.com/weaveworks/scope/test" ) -const ( - release = "release" - version = "version" - network = "192.168.0.0/16" - hostID = "hostid" - now = "now" - hostname = "hostname" - load = "0.59 0.36 0.29" - uptime = "278h55m43s" - kernel = "release version" -) - func TestReporter(t *testing.T) { + var ( + release = "release" + version = "version" + network = "192.168.0.0/16" + hostID = "hostid" + now = "now" + hostname = "hostname" + load = "0.59 0.36 0.29" + uptime = "278h55m43s" + kernel = "release version" + _, ipnet, _ = net.ParseCIDR(network) + localNets = report.Networks([]*net.IPNet{ipnet}) + ) + var ( oldGetKernelVersion = host.GetKernelVersion oldGetLoad = host.GetLoad oldGetUptime = host.GetUptime - oldInterfaceAddrs = host.InterfaceAddrs oldNow = host.Now ) defer func() { host.GetKernelVersion = oldGetKernelVersion host.GetLoad = oldGetLoad host.GetUptime = oldGetUptime - host.InterfaceAddrs = oldInterfaceAddrs host.Now = oldNow }() host.GetKernelVersion = func() (string, error) { return release + " " + version, nil } host.GetLoad = func() string { return load } host.GetUptime = func() (time.Duration, error) { return time.ParseDuration(uptime) } host.Now = func() string { return now } - host.InterfaceAddrs = func() ([]net.Addr, error) { _, ipnet, _ := net.ParseCIDR(network); return []net.Addr{ipnet}, nil } want := report.MakeReport() want.Host.NodeMetadatas[report.MakeHostNodeID(hostID)] = report.MakeNodeMetadataWith(map[string]string{ @@ -55,8 +54,7 @@ func TestReporter(t *testing.T) { host.Uptime: uptime, host.KernelVersion: kernel, }) - r := host.NewReporter(hostID, hostname) - have, _ := r.Report() + have, _ := host.NewReporter(hostID, hostname, localNets).Report() if !reflect.DeepEqual(want, have) { t.Errorf("%s", test.Diff(want, have)) } diff --git a/probe/main.go b/probe/main.go index 0b3a8c63a..29f7626d1 100644 --- a/probe/main.go +++ b/probe/main.go @@ -79,11 +79,23 @@ func main() { } defer publisher.Close() + addrs, err := net.InterfaceAddrs() + if err != nil { + log.Fatal(err) + } + localNets := report.Networks{} + for _, addr := range addrs { + // Not all addrs are IPNets. + if ipNet, ok := addr.(*net.IPNet); ok { + localNets = append(localNets, ipNet) + } + } + var ( hostName = hostname() hostID = hostName // TODO: we should sanitize the hostname taggers = []Tagger{newTopologyTagger(), host.NewTagger(hostID)} - reporters = []Reporter{host.NewReporter(hostID, hostName), endpoint.NewReporter(hostID, hostName, *spyProcs)} + reporters = []Reporter{host.NewReporter(hostID, hostName, localNets), endpoint.NewReporter(hostID, hostName, *spyProcs)} processCache *process.CachingWalker ) @@ -122,7 +134,7 @@ func main() { continue } log.Printf("capturing packets on %s", iface) - reporters = append(reporters, sniff.New(hostID, source, *captureOn, *captureOff)) + reporters = append(reporters, sniff.New(hostID, localNets, source, *captureOn, *captureOff)) } } diff --git a/probe/sniff/sniffer.go b/probe/sniff/sniffer.go index e0a13feef..91dff56c5 100644 --- a/probe/sniff/sniffer.go +++ b/probe/sniff/sniffer.go @@ -3,6 +3,7 @@ package sniff import ( "io" "log" + "net" "strconv" "sync/atomic" "time" @@ -15,27 +16,29 @@ import ( // Sniffer is a packet-sniffing reporter. type Sniffer struct { - hostID string - reports chan chan report.Report - parser *gopacket.DecodingLayerParser - decoded []gopacket.LayerType - eth layers.Ethernet - ip4 layers.IPv4 - ip6 layers.IPv6 - tcp layers.TCP - udp layers.UDP - icmp4 layers.ICMPv4 - icmp6 layers.ICMPv6 + hostID string + localNets report.Networks + reports chan chan report.Report + parser *gopacket.DecodingLayerParser + decoded []gopacket.LayerType + eth layers.Ethernet + ip4 layers.IPv4 + ip6 layers.IPv6 + tcp layers.TCP + udp layers.UDP + icmp4 layers.ICMPv4 + icmp6 layers.ICMPv6 } // New returns a new sniffing reporter that samples traffic by turning its // packet capture facilities on and off. Note that the on and off durations // represent a way to bound CPU burn. Effective sample rate needs to be // calculated as (packets decoded / packets observed). -func New(hostID string, src gopacket.ZeroCopyPacketDataSource, on, off time.Duration) *Sniffer { +func New(hostID string, localNets report.Networks, src gopacket.ZeroCopyPacketDataSource, on, off time.Duration) *Sniffer { s := &Sniffer{ - hostID: hostID, - reports: make(chan chan report.Report), + hostID: hostID, + localNets: localNets, + reports: make(chan chan report.Report), } s.parser = gopacket.NewDecodingLayerParser( layers.LayerTypeEthernet, @@ -119,8 +122,11 @@ func interpolateCounts(r report.Report) { if emd.PacketCount != nil { *emd.PacketCount = uint64(float64(*emd.PacketCount) * factor) } - if emd.ByteCount != nil { - *emd.ByteCount = uint64(float64(*emd.ByteCount) * factor) + if emd.EgressByteCount != nil { + *emd.EgressByteCount = uint64(float64(*emd.EgressByteCount) * factor) + } + if emd.IngressByteCount != nil { + *emd.IngressByteCount = uint64(float64(*emd.IngressByteCount) * factor) } } } @@ -204,54 +210,104 @@ func (s *Sniffer) read(src gopacket.ZeroCopyPacketDataSource, dst chan Packet, p } // Merge puts the packet into the report. +// +// Note that, for the moment, we encode bidirectional traffic as ingress and +// egress traffic on a single edge whose src is local and dst is remote. That +// is, if we see a packet from the remote addr 9.8.7.6 to the local addr +// 1.2.3.4, we apply it as *ingress* on the edge (1.2.3.4 -> 9.8.7.6). func (s *Sniffer) Merge(p Packet, rpt report.Report) { - // With a src and dst IP, we can add to the address topology. - if p.SrcIP != "" && p.DstIP != "" { + if p.SrcIP == "" || p.DstIP == "" { + return + } + + // One end of the traffic has to be local. Otherwise, we don't know how to + // construct the edge. + // + // If we need to get around this limitation, we may be able to change the + // semantics of the report, and allow the src side of edges to be from + // anywhere. But that will have ramifications throughout Scope (read: it + // may violate implicit invariants) and needs to be thought through. + var ( + srcLocal = s.localNets.Contains(net.ParseIP(p.SrcIP)) + dstLocal = s.localNets.Contains(net.ParseIP(p.DstIP)) + localIP string + remoteIP string + egress bool + ) + switch { + case srcLocal && !dstLocal: + localIP, remoteIP, egress = p.SrcIP, p.DstIP, true + case !srcLocal && dstLocal: + localIP, remoteIP, egress = p.DstIP, p.SrcIP, false + case srcLocal && dstLocal: + localIP, remoteIP, egress = p.SrcIP, p.DstIP, true // loopback + case !srcLocal && !dstLocal: + log.Printf("sniffer ignoring remote-to-remote (%s -> %s) traffic", p.SrcIP, p.DstIP) + return + } + + // For sure, we can add to the address topology. + { var ( - srcNodeID = report.MakeAddressNodeID(s.hostID, p.SrcIP) - dstNodeID = report.MakeAddressNodeID(s.hostID, p.DstIP) + srcNodeID = report.MakeAddressNodeID(s.hostID, localIP) + dstNodeID = report.MakeAddressNodeID(s.hostID, remoteIP) edgeID = report.MakeEdgeID(srcNodeID, dstNodeID) srcAdjacencyID = report.MakeAdjacencyID(srcNodeID) ) + rpt.Address.NodeMetadatas[srcNodeID] = report.MakeNodeMetadata() - rpt.Address.NodeMetadatas[dstNodeID] = report.MakeNodeMetadata() emd := rpt.Address.EdgeMetadatas[edgeID] if emd.PacketCount == nil { emd.PacketCount = new(uint64) } *emd.PacketCount++ - if emd.ByteCount == nil { - emd.ByteCount = new(uint64) - } - *emd.ByteCount += uint64(p.Network) - rpt.Address.EdgeMetadatas[edgeID] = emd + if egress { + if emd.EgressByteCount == nil { + emd.EgressByteCount = new(uint64) + } + *emd.EgressByteCount += uint64(p.Network) + } else { + if emd.IngressByteCount == nil { + emd.IngressByteCount = new(uint64) + } + *emd.IngressByteCount += uint64(p.Network) + } + + rpt.Address.EdgeMetadatas[edgeID] = emd rpt.Address.Adjacency[srcAdjacencyID] = rpt.Address.Adjacency[srcAdjacencyID].Add(dstNodeID) } - // With a src and dst IP and port, we can add to the endpoints. - if p.SrcIP != "" && p.DstIP != "" && p.SrcPort != "" && p.DstPort != "" { + // If we have ports, we can add to the endpoint topology, too. + if p.SrcPort != "" && p.DstPort != "" { var ( - srcNodeID = report.MakeEndpointNodeID(s.hostID, p.SrcIP, p.SrcPort) - dstNodeID = report.MakeEndpointNodeID(s.hostID, p.DstIP, p.DstPort) + srcNodeID = report.MakeEndpointNodeID(s.hostID, localIP, p.SrcPort) + dstNodeID = report.MakeEndpointNodeID(s.hostID, remoteIP, p.DstPort) edgeID = report.MakeEdgeID(srcNodeID, dstNodeID) srcAdjacencyID = report.MakeAdjacencyID(srcNodeID) ) rpt.Endpoint.NodeMetadatas[srcNodeID] = report.MakeNodeMetadata() - rpt.Endpoint.NodeMetadatas[dstNodeID] = report.MakeNodeMetadata() emd := rpt.Endpoint.EdgeMetadatas[edgeID] if emd.PacketCount == nil { emd.PacketCount = new(uint64) } *emd.PacketCount++ - if emd.ByteCount == nil { - emd.ByteCount = new(uint64) - } - *emd.ByteCount += uint64(p.Transport) - rpt.Endpoint.EdgeMetadatas[edgeID] = emd + if egress { + if emd.EgressByteCount == nil { + emd.EgressByteCount = new(uint64) + } + *emd.EgressByteCount += uint64(p.Transport) + } else { + if emd.IngressByteCount == nil { + emd.IngressByteCount = new(uint64) + } + *emd.IngressByteCount += uint64(p.Transport) + } + + rpt.Endpoint.EdgeMetadatas[edgeID] = emd rpt.Endpoint.Adjacency[srcAdjacencyID] = rpt.Endpoint.Adjacency[srcAdjacencyID].Add(dstNodeID) } } diff --git a/probe/sniff/sniffer_internal_test.go b/probe/sniff/sniffer_internal_test.go index ed6d559a2..fc04b2245 100644 --- a/probe/sniff/sniffer_internal_test.go +++ b/probe/sniff/sniffer_internal_test.go @@ -22,8 +22,9 @@ func TestInterpolateCounts(t *testing.T) { r.Sampling.Count = samplingCount r.Sampling.Total = samplingTotal r.Endpoint.EdgeMetadatas[edgeID] = report.EdgeMetadata{ - PacketCount: newu64(packetCount), - ByteCount: newu64(byteCount), + PacketCount: newu64(packetCount), + IngressByteCount: newu64(byteCount), + EgressByteCount: newu64(byteCount), } interpolateCounts(r) @@ -37,7 +38,10 @@ func TestInterpolateCounts(t *testing.T) { if want, have := apply(packetCount), (*emd.PacketCount); want != have { t.Errorf("want %d packets, have %d", want, have) } - if want, have := apply(byteCount), (*emd.ByteCount); want != have { + if want, have := apply(byteCount), (*emd.EgressByteCount); want != have { + t.Errorf("want %d bytes, have %d", want, have) + } + if want, have := apply(byteCount), (*emd.IngressByteCount); want != have { t.Errorf("want %d bytes, have %d", want, have) } } diff --git a/probe/sniff/sniffer_test.go b/probe/sniff/sniffer_test.go index 3a686e382..6187924bc 100644 --- a/probe/sniff/sniffer_test.go +++ b/probe/sniff/sniffer_test.go @@ -2,6 +2,7 @@ package sniff_test import ( "io" + "net" "reflect" "sync" "testing" @@ -20,7 +21,7 @@ func TestSnifferShutdown(t *testing.T) { src = newMockSource([]byte{}, nil) on = time.Millisecond off = time.Millisecond - s = sniff.New(hostID, src, on, off) + s = sniff.New(hostID, report.Networks{}, src, on, off) ) // Stopping the source should terminate the sniffer. @@ -53,8 +54,11 @@ func TestMerge(t *testing.T) { Network: 512, Transport: 256, } + + _, ipnet, _ = net.ParseCIDR(p.SrcIP + "/24") // ;) + localNets = report.Networks([]*net.IPNet{ipnet}) ) - sniff.New(hostID, src, on, off).Merge(p, rpt) + sniff.New(hostID, localNets, src, on, off).Merge(p, rpt) var ( srcEndpointNodeID = report.MakeEndpointNodeID(hostID, p.SrcIP, p.SrcPort) @@ -68,13 +72,12 @@ func TestMerge(t *testing.T) { }, EdgeMetadatas: report.EdgeMetadatas{ report.MakeEdgeID(srcEndpointNodeID, dstEndpointNodeID): report.EdgeMetadata{ - PacketCount: newu64(1), - ByteCount: newu64(256), + PacketCount: newu64(1), + EgressByteCount: newu64(256), }, }, NodeMetadatas: report.NodeMetadatas{ srcEndpointNodeID: report.MakeNodeMetadata(), - dstEndpointNodeID: report.MakeNodeMetadata(), }, }), rpt.Endpoint; !reflect.DeepEqual(want, have) { t.Errorf("%s", test.Diff(want, have)) @@ -92,13 +95,12 @@ func TestMerge(t *testing.T) { }, EdgeMetadatas: report.EdgeMetadatas{ report.MakeEdgeID(srcAddressNodeID, dstAddressNodeID): report.EdgeMetadata{ - PacketCount: newu64(1), - ByteCount: newu64(512), + PacketCount: newu64(1), + EgressByteCount: newu64(512), }, }, NodeMetadatas: report.NodeMetadatas{ srcAddressNodeID: report.MakeNodeMetadata(), - dstAddressNodeID: report.MakeNodeMetadata(), }, }), rpt.Address; !reflect.DeepEqual(want, have) { t.Errorf("%s", test.Diff(want, have)) diff --git a/render/detailed_node.go b/render/detailed_node.go index c1f0def81..fc05a4ad0 100644 --- a/render/detailed_node.go +++ b/render/detailed_node.go @@ -67,8 +67,11 @@ func MakeDetailedNode(r report.Report, n RenderableNode) DetailedNode { if n.EdgeMetadata.PacketCount != nil { rows = append(rows, Row{"Packets", strconv.FormatUint(*n.EdgeMetadata.PacketCount, 10), ""}) } - if n.EdgeMetadata.ByteCount != nil { - rows = append(rows, Row{"Bytes", strconv.FormatUint(*n.EdgeMetadata.ByteCount, 10), ""}) + if n.EdgeMetadata.EgressByteCount != nil { + rows = append(rows, Row{"Egress bytes", strconv.FormatUint(*n.EdgeMetadata.EgressByteCount, 10), ""}) // TODO rate + } + if n.EdgeMetadata.IngressByteCount != nil { + rows = append(rows, Row{"Ingress bytes", strconv.FormatUint(*n.EdgeMetadata.IngressByteCount, 10), ""}) // TODO rate } if len(rows) > 0 { tables = append(tables, Table{"Connections", true, connectionsRank, rows}) diff --git a/render/detailed_node_test.go b/render/detailed_node_test.go index 0b3fe9994..606f05556 100644 --- a/render/detailed_node_test.go +++ b/render/detailed_node_test.go @@ -67,7 +67,7 @@ func TestMakeDetailedNode(t *testing.T) { Rank: 100, Rows: []render.Row{ {"Packets", "150", ""}, - {"Bytes", "1500", ""}, + {"Egress bytes", "1500", ""}, }, }, { diff --git a/render/expected/expected.go b/render/expected/expected.go index 7e3b5fd74..5d4ebc634 100644 --- a/render/expected/expected.go +++ b/render/expected/expected.go @@ -55,8 +55,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(100), - ByteCount: newu64(10), + PacketCount: newu64(10), + EgressByteCount: newu64(100), }, }, ClientProcess2ID: { @@ -73,8 +73,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(200), - ByteCount: newu64(20), + PacketCount: newu64(20), + EgressByteCount: newu64(200), }, }, ServerProcessID: { @@ -97,8 +97,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(150), - ByteCount: newu64(1500), + PacketCount: newu64(150), + EgressByteCount: newu64(1500), }, }, nonContainerProcessID: { @@ -137,8 +137,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(300), - ByteCount: newu64(30), + PacketCount: newu64(30), + EgressByteCount: newu64(300), }, }, "apache": { @@ -160,8 +160,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(150), - ByteCount: newu64(1500), + PacketCount: newu64(150), + EgressByteCount: newu64(1500), }, }, "bash": { @@ -200,8 +200,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(300), - ByteCount: newu64(30), + PacketCount: newu64(30), + EgressByteCount: newu64(300), }, }, test.ServerContainerID: { @@ -219,8 +219,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(150), - ByteCount: newu64(1500), + PacketCount: newu64(150), + EgressByteCount: newu64(1500), }, }, uncontainedServerID: { @@ -258,8 +258,8 @@ var ( ), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(300), - ByteCount: newu64(30), + PacketCount: newu64(30), + EgressByteCount: newu64(300), }, }, test.ServerContainerImageName: { @@ -277,8 +277,8 @@ var ( test.ServerHostNodeID), NodeMetadata: report.MakeNodeMetadata(), EdgeMetadata: report.EdgeMetadata{ - PacketCount: newu64(150), - ByteCount: newu64(1500), + PacketCount: newu64(150), + EgressByteCount: newu64(1500), }, }, uncontainedServerID: { diff --git a/render/render_test.go b/render/render_test.go index 954d9df34..0c50cbb7d 100644 --- a/render/render_test.go +++ b/render/render_test.go @@ -118,8 +118,8 @@ func TestMapEdge(t *testing.T) { ">bar": report.MakeIDList("foo"), }, EdgeMetadatas: report.EdgeMetadatas{ - "foo|bar": report.EdgeMetadata{PacketCount: newu64(1), ByteCount: newu64(2)}, - "bar|foo": report.EdgeMetadata{PacketCount: newu64(3), ByteCount: newu64(4)}, + "foo|bar": report.EdgeMetadata{PacketCount: newu64(1), EgressByteCount: newu64(2)}, + "bar|foo": report.EdgeMetadata{PacketCount: newu64(3), EgressByteCount: newu64(4)}, }, } } @@ -140,8 +140,8 @@ func TestMapEdge(t *testing.T) { } if want, have := (report.EdgeMetadata{ - PacketCount: newu64(1), - ByteCount: newu64(2), + PacketCount: newu64(1), + EgressByteCount: newu64(2), }), mapper.EdgeMetadata(report.MakeReport(), "_foo", "_bar"); !reflect.DeepEqual(want, have) { t.Errorf("want %+v, have %+v", want, have) } diff --git a/report/merge.go b/report/merge.go index 8ebd1eab2..12c053105 100644 --- a/report/merge.go +++ b/report/merge.go @@ -64,7 +64,8 @@ func (e *EdgeMetadatas) Merge(other EdgeMetadatas) { // should represent the same edge on different times. func (m *EdgeMetadata) Merge(other EdgeMetadata) { m.PacketCount = merge(m.PacketCount, other.PacketCount, sum) - m.ByteCount = merge(m.ByteCount, other.ByteCount, sum) + m.EgressByteCount = merge(m.EgressByteCount, other.EgressByteCount, sum) + m.IngressByteCount = merge(m.IngressByteCount, other.IngressByteCount, sum) m.MaxConnCountTCP = merge(m.MaxConnCountTCP, other.MaxConnCountTCP, max) } @@ -72,7 +73,8 @@ func (m *EdgeMetadata) Merge(other EdgeMetadata) { // they should represent different edges at the same time. func (m *EdgeMetadata) Flatten(other EdgeMetadata) { m.PacketCount = merge(m.PacketCount, other.PacketCount, sum) - m.ByteCount = merge(m.ByteCount, other.ByteCount, sum) + m.EgressByteCount = merge(m.EgressByteCount, other.EgressByteCount, sum) + m.IngressByteCount = merge(m.IngressByteCount, other.IngressByteCount, sum) // Note that summing of two maximums doesn't always give us the true // maximum. But it's a best effort. m.MaxConnCountTCP = merge(m.MaxConnCountTCP, other.MaxConnCountTCP, sum) diff --git a/report/merge_test.go b/report/merge_test.go index 5f3a75beb..b87ea176e 100644 --- a/report/merge_test.go +++ b/report/merge_test.go @@ -112,15 +112,13 @@ func TestMergeEdgeMetadatas(t *testing.T) { a: report.EdgeMetadatas{}, b: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(12), - ByteCount: newu64(0), + PacketCount: newu64(1), MaxConnCountTCP: newu64(2), }, }, want: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(12), - ByteCount: newu64(0), + PacketCount: newu64(1), MaxConnCountTCP: newu64(2), }, }, @@ -128,15 +126,15 @@ func TestMergeEdgeMetadatas(t *testing.T) { "Empty b": { a: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(12), - ByteCount: newu64(0), + PacketCount: newu64(12), + EgressByteCount: newu64(999), }, }, b: report.EdgeMetadatas{}, want: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(12), - ByteCount: newu64(0), + PacketCount: newu64(12), + EgressByteCount: newu64(999), }, }, }, @@ -144,26 +142,26 @@ func TestMergeEdgeMetadatas(t *testing.T) { a: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ PacketCount: newu64(12), - ByteCount: newu64(0), + EgressByteCount: newu64(500), MaxConnCountTCP: newu64(4), }, }, b: report.EdgeMetadatas{ "hostQ|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ PacketCount: newu64(1), - ByteCount: newu64(2), + EgressByteCount: newu64(2), MaxConnCountTCP: newu64(6), }, }, want: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ PacketCount: newu64(12), - ByteCount: newu64(0), + EgressByteCount: newu64(500), MaxConnCountTCP: newu64(4), }, "hostQ|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ PacketCount: newu64(1), - ByteCount: newu64(2), + EgressByteCount: newu64(2), MaxConnCountTCP: newu64(6), }, }, @@ -172,22 +170,24 @@ func TestMergeEdgeMetadatas(t *testing.T) { a: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ PacketCount: newu64(12), - ByteCount: newu64(0), + EgressByteCount: newu64(1000), MaxConnCountTCP: newu64(7), }, }, b: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(1), - ByteCount: newu64(2), - MaxConnCountTCP: newu64(9), + PacketCount: newu64(1), + IngressByteCount: newu64(123), + EgressByteCount: newu64(2), + MaxConnCountTCP: newu64(9), }, }, want: report.EdgeMetadatas{ "hostA|:192.168.1.1:12345|:192.168.1.2:80": report.EdgeMetadata{ - PacketCount: newu64(13), - ByteCount: newu64(2), - MaxConnCountTCP: newu64(9), + PacketCount: newu64(13), + IngressByteCount: newu64(123), + EgressByteCount: newu64(1002), + MaxConnCountTCP: newu64(9), }, }, }, diff --git a/report/networks.go b/report/networks.go index 4fe5f6fe4..a1aa9e6e4 100644 --- a/report/networks.go +++ b/report/networks.go @@ -12,12 +12,11 @@ type Interface interface { Addrs() ([]net.Addr, error) } -// Variables exposed for testing +// Variables exposed for testing. +// TODO this design is broken, make it consistent with probe networks. var ( LocalNetworks = Networks{} - InterfaceByNameStub = func(name string) (Interface, error) { - return net.InterfaceByName(name) - } + InterfaceByNameStub = func(name string) (Interface, error) { return net.InterfaceByName(name) } ) // Contains returns true if IP is in Networks. diff --git a/report/topology.go b/report/topology.go index 9b9a5fed2..0327f35ad 100644 --- a/report/topology.go +++ b/report/topology.go @@ -31,9 +31,10 @@ type NodeMetadatas map[string]NodeMetadata // EdgeMetadata describes a superset of the metadata that probes can possibly // collect about a directed edge between two nodes in any topology. type EdgeMetadata struct { - PacketCount *uint64 `json:"packet_count,omitempty"` - ByteCount *uint64 `json:"byte_count,omitempty"` - MaxConnCountTCP *uint64 `json:"max_conn_count_tcp,omitempty"` + PacketCount *uint64 `json:"packet_count,omitempty"` + EgressByteCount *uint64 `json:"ingress_byte_count,omitempty"` + IngressByteCount *uint64 `json:"egress_byte_count,omitempty"` + MaxConnCountTCP *uint64 `json:"max_conn_count_tcp,omitempty"` } // NodeMetadata describes a superset of the metadata that probes can collect diff --git a/test/report_fixture.go b/test/report_fixture.go index b48fcd8dc..c949a2ae1 100644 --- a/test/report_fixture.go +++ b/test/report_fixture.go @@ -108,33 +108,33 @@ var ( }, EdgeMetadatas: report.EdgeMetadatas{ report.MakeEdgeID(Client54001NodeID, Server80NodeID): report.EdgeMetadata{ - PacketCount: newu64(100), - ByteCount: newu64(10), + PacketCount: newu64(10), + EgressByteCount: newu64(100), }, report.MakeEdgeID(Client54002NodeID, Server80NodeID): report.EdgeMetadata{ - PacketCount: newu64(200), - ByteCount: newu64(20), + PacketCount: newu64(20), + EgressByteCount: newu64(200), }, report.MakeEdgeID(Server80NodeID, Client54001NodeID): report.EdgeMetadata{ - PacketCount: newu64(10), - ByteCount: newu64(100), + PacketCount: newu64(10), + EgressByteCount: newu64(100), }, report.MakeEdgeID(Server80NodeID, Client54002NodeID): report.EdgeMetadata{ - PacketCount: newu64(20), - ByteCount: newu64(200), + PacketCount: newu64(20), + EgressByteCount: newu64(200), }, report.MakeEdgeID(Server80NodeID, UnknownClient1NodeID): report.EdgeMetadata{ - PacketCount: newu64(30), - ByteCount: newu64(300), + PacketCount: newu64(30), + EgressByteCount: newu64(300), }, report.MakeEdgeID(Server80NodeID, UnknownClient2NodeID): report.EdgeMetadata{ - PacketCount: newu64(40), - ByteCount: newu64(400), + PacketCount: newu64(40), + EgressByteCount: newu64(400), }, report.MakeEdgeID(Server80NodeID, UnknownClient3NodeID): report.EdgeMetadata{ - PacketCount: newu64(50), - ByteCount: newu64(500), + PacketCount: newu64(50), + EgressByteCount: newu64(500), }, }, },