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), }, }, },