From 15e25edc40c36cd6760867aaddaae138cbddd918 Mon Sep 17 00:00:00 2001 From: Alvaro Saurin Date: Thu, 27 Aug 2015 14:37:53 +0200 Subject: [PATCH 1/3] New asynchronous, caching DNS resolver for reverse resolutions Add nodes for the remote side of connections iff we have a DNS reverse resolution for the IP. Unit test for the resolver --- probe/endpoint/reporter.go | 20 ++++++ probe/endpoint/reporter_test.go | 2 +- probe/endpoint/resolver.go | 85 ++++++++++++++++++++++++ probe/endpoint/resolver_internal_test.go | 41 ++++++++++++ render/detailed_node.go | 11 +-- report/topology.go | 6 ++ 6 files changed, 160 insertions(+), 5 deletions(-) create mode 100644 probe/endpoint/resolver.go create mode 100644 probe/endpoint/resolver_internal_test.go diff --git a/probe/endpoint/reporter.go b/probe/endpoint/reporter.go index bd71224b7..049659b13 100644 --- a/probe/endpoint/reporter.go +++ b/probe/endpoint/reporter.go @@ -26,6 +26,7 @@ type Reporter struct { includeNAT bool conntracker *Conntracker natmapper *natmapper + revResolver *reverseResolver } // SpyDuration is an exported prometheus metric @@ -64,12 +65,16 @@ func NewReporter(hostID, hostName string, includeProcesses bool, useConntrack bo log.Printf("Failed to start natMapper: %v", err) } } + + revRes := newReverseResolver() + return &Reporter{ hostID: hostID, hostName: hostName, includeProcesses: includeProcesses, conntracker: conntracker, natmapper: natmapper, + revResolver: revRes, } } @@ -81,6 +86,7 @@ func (r *Reporter) Stop() { if r.natmapper != nil { r.natmapper.Stop() } + r.revResolver.Stop() } // Report implements Reporter. @@ -145,6 +151,13 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin }) ) + // in case we have a reverse resolution for the IP, we can use it for the name... + if revRemoteName, err := r.revResolver.Get(remoteAddr, false); err == nil { + remoteNode = remoteNode.AddMetadata(map[string]string{ + "name": revRemoteName, + }) + } + if localIsClient { // New nodes are merged into the report so we don't need to do any counting here; the merge does it for us. localNode = localNode.WithEdge(remoteAddressNodeID, report.EdgeMetadata{ @@ -177,6 +190,13 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin }) ) + // in case we have a reverse resolution for the IP, we can use it for the name... + if revRemoteName, err := r.revResolver.Get(remoteAddr, false); err == nil { + remoteNode = remoteNode.AddMetadata(map[string]string{ + "name": revRemoteName, + }) + } + if localIsClient { // New nodes are merged into the report so we don't need to do any counting here; the merge does it for us. localNode = localNode.WithEdge(remoteEndpointNodeID, report.EdgeMetadata{ diff --git a/probe/endpoint/reporter_test.go b/probe/endpoint/reporter_test.go index d3e3fb107..c7ec41995 100644 --- a/probe/endpoint/reporter_test.go +++ b/probe/endpoint/reporter_test.go @@ -54,7 +54,7 @@ var ( LocalAddress: fixLocalAddress, LocalPort: fixLocalPort, RemoteAddress: fixRemoteAddress, - RemotePort: fixRemotePortB, + RemotePort: fixRemotePort, Proc: procspy.Proc{ PID: fixProcessPID, Name: fixProcessName, diff --git a/probe/endpoint/resolver.go b/probe/endpoint/resolver.go new file mode 100644 index 000000000..c27bef9f4 --- /dev/null +++ b/probe/endpoint/resolver.go @@ -0,0 +1,85 @@ +package endpoint + +import ( + "net" + "strings" + "time" + + "github.com/bluele/gcache" +) + +const ( + rAddrCacheLen = 500 // Default cache length + rAddrBacklog = 1000 + rAddrCacheExpiration = 30 * time.Minute +) + +type revResFunc func(addr string) (names []string, err error) + +type revResRequest struct { + address string + done chan struct{} +} + +// ReverseResolver is a caching, reverse resolver +type reverseResolver struct { + addresses chan revResRequest + cache gcache.Cache + resolver revResFunc +} + +// NewReverseResolver starts a new reverse resolver that +// performs reverse resolutions and caches the result. +func newReverseResolver() *reverseResolver { + r := reverseResolver{ + addresses: make(chan revResRequest, rAddrBacklog), + cache: gcache.New(rAddrCacheLen).LRU().Expiration(rAddrCacheExpiration).Build(), + resolver: net.LookupAddr, + } + go r.loop() + return &r +} + +// Get the reverse resolution for an IP address if already in the cache, +// a gcache.NotFoundKeyError error otherwise. +// Note: it returns one of the possible names that can be obtained for that IP. +func (r *reverseResolver) Get(address string, wait bool) (string, error) { + val, err := r.cache.Get(address) + if err == nil { + return val.(string), nil + } + if err == gcache.NotFoundKeyError { + request := revResRequest{address: address, done: make(chan struct{})} + // we trigger a asynchronous reverse resolution when not cached + select { + case r.addresses <- request: + if wait { + <-request.done + } + default: + } + } + return "", err +} + +func (r *reverseResolver) loop() { + throttle := time.Tick(time.Second / 10) + for request := range r.addresses { + <-throttle // rate limit our DNS resolutions + // and check if the answer is already in the cache + if _, err := r.cache.Get(request.address); err == nil { + continue + } + names, err := r.resolver(request.address) + if err == nil && len(names) > 0 { + name := strings.TrimRight(names[0], ".") + r.cache.Set(request.address, name) + } + close(request.done) + } +} + +// Stop the async reverse resolver +func (r *reverseResolver) Stop() { + close(r.addresses) +} diff --git a/probe/endpoint/resolver_internal_test.go b/probe/endpoint/resolver_internal_test.go new file mode 100644 index 000000000..e9a2d264e --- /dev/null +++ b/probe/endpoint/resolver_internal_test.go @@ -0,0 +1,41 @@ +package endpoint + +import ( + "errors" + "testing" +) + +func TestReverseResolver(t *testing.T) { + tests := map[string]string{ + "8.8.8.8": "google-public-dns-a.google.com", + "8.8.4.4": "google-public-dns-b.google.com", + } + + revRes := newReverseResolver() + + // use a mocked resolver function + revRes.resolver = func(addr string) (names []string, err error) { + if name, ok := tests[addr]; ok { + return []string{name}, nil + } + return []string{}, errors.New("invalid IP") + } + + // first time: no names are returned for our reverse resolutions + for ip := range tests { + if have, err := revRes.Get(ip, true); have != "" || err == nil { + t.Errorf("we didn't get an error, or the cache was not empty, when trying to resolve '%q'", ip) + } + } + + // so, if we check again these IPs, we should have the names now + for ip, want := range tests { + have, err := revRes.Get(ip, true) + if err != nil { + t.Errorf("%s: %v", ip, err) + } + if want != have { + t.Errorf("%s: want %q, have %q", ip, want, have) + } + } +} diff --git a/render/detailed_node.go b/render/detailed_node.go index 7459aea7b..6e40e3532 100644 --- a/render/detailed_node.go +++ b/render/detailed_node.go @@ -205,8 +205,11 @@ func OriginTable(r report.Report, originID string, addHostTags bool, addContaine func connectionDetailsRows(topology report.Topology, originID string) []Row { rows := []Row{} - labeler := func(nodeID string) (string, bool) { + labeler := func(nodeID string, meta map[string]string) (string, bool) { if _, addr, port, ok := report.ParseEndpointNodeID(nodeID); ok { + if name, ok := meta["name"]; ok { + return fmt.Sprintf("%s:%s", name, port), true + } return fmt.Sprintf("%s:%s", addr, port), true } if _, addr, ok := report.ParseAddressNodeID(nodeID); ok { @@ -214,13 +217,13 @@ func connectionDetailsRows(topology report.Topology, originID string) []Row { } return "", false } - local, ok := labeler(originID) + local, ok := labeler(originID, topology.Nodes[originID].Metadata) if !ok { return rows } // Firstly, collection outgoing connections from this node. for _, serverNodeID := range topology.Nodes[originID].Adjacency { - remote, ok := labeler(serverNodeID) + remote, ok := labeler(serverNodeID, topology.Nodes[serverNodeID].Metadata) if !ok { continue } @@ -239,7 +242,7 @@ func connectionDetailsRows(topology report.Topology, originID string) []Row { if !serverNodeIDs.Contains(originID) { continue } - remote, ok := labeler(clientNodeID) + remote, ok := labeler(clientNodeID, clientNode.Metadata) if !ok { continue } diff --git a/report/topology.go b/report/topology.go index a938b7095..bf237d996 100644 --- a/report/topology.go +++ b/report/topology.go @@ -102,6 +102,12 @@ func (n Node) WithMetadata(m map[string]string) Node { return result } +// AddMetadata returns a fresh copy of n, with Metadata set to the merge of n and the metadata provided +func (n Node) AddMetadata(m map[string]string) Node { + additional := MakeNodeWith(m) + return n.Merge(additional) +} + // WithCounters returns a fresh copy of n, with Counters set to c func (n Node) WithCounters(c map[string]int) Node { result := n.Copy() From ad6702a196b301834868f2d47c5a4352fe39c170 Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Sat, 5 Sep 2015 19:39:22 +0000 Subject: [PATCH 2/3] Some review feedback. --- probe/endpoint/reporter.go | 11 +++--- probe/endpoint/resolver.go | 43 ++++++++++-------------- probe/endpoint/resolver_internal_test.go | 41 ---------------------- probe/endpoint/resolver_test.go | 38 +++++++++++++++++++++ 4 files changed, 59 insertions(+), 74 deletions(-) delete mode 100644 probe/endpoint/resolver_internal_test.go create mode 100644 probe/endpoint/resolver_test.go diff --git a/probe/endpoint/reporter.go b/probe/endpoint/reporter.go index 049659b13..8685fa932 100644 --- a/probe/endpoint/reporter.go +++ b/probe/endpoint/reporter.go @@ -26,7 +26,7 @@ type Reporter struct { includeNAT bool conntracker *Conntracker natmapper *natmapper - revResolver *reverseResolver + revResolver *ReverseResolver } // SpyDuration is an exported prometheus metric @@ -65,16 +65,13 @@ func NewReporter(hostID, hostName string, includeProcesses bool, useConntrack bo log.Printf("Failed to start natMapper: %v", err) } } - - revRes := newReverseResolver() - return &Reporter{ hostID: hostID, hostName: hostName, includeProcesses: includeProcesses, conntracker: conntracker, natmapper: natmapper, - revResolver: revRes, + revResolver: NewReverseResolver(), } } @@ -152,7 +149,7 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin ) // in case we have a reverse resolution for the IP, we can use it for the name... - if revRemoteName, err := r.revResolver.Get(remoteAddr, false); err == nil { + if revRemoteName, err := r.revResolver.Get(remoteAddr); err == nil { remoteNode = remoteNode.AddMetadata(map[string]string{ "name": revRemoteName, }) @@ -191,7 +188,7 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin ) // in case we have a reverse resolution for the IP, we can use it for the name... - if revRemoteName, err := r.revResolver.Get(remoteAddr, false); err == nil { + if revRemoteName, err := r.revResolver.Get(remoteAddr); err == nil { remoteNode = remoteNode.AddMetadata(map[string]string{ "name": revRemoteName, }) diff --git a/probe/endpoint/resolver.go b/probe/endpoint/resolver.go index c27bef9f4..210e93832 100644 --- a/probe/endpoint/resolver.go +++ b/probe/endpoint/resolver.go @@ -16,25 +16,22 @@ const ( type revResFunc func(addr string) (names []string, err error) -type revResRequest struct { - address string - done chan struct{} -} - // ReverseResolver is a caching, reverse resolver -type reverseResolver struct { - addresses chan revResRequest +type ReverseResolver struct { + addresses chan string cache gcache.Cache - resolver revResFunc + Throttle <-chan time.Time // Made public for mocking + Resolver revResFunc } // NewReverseResolver starts a new reverse resolver that // performs reverse resolutions and caches the result. -func newReverseResolver() *reverseResolver { - r := reverseResolver{ - addresses: make(chan revResRequest, rAddrBacklog), +func NewReverseResolver() *ReverseResolver { + r := ReverseResolver{ + addresses: make(chan string, rAddrBacklog), cache: gcache.New(rAddrCacheLen).LRU().Expiration(rAddrCacheExpiration).Build(), - resolver: net.LookupAddr, + Throttle: time.Tick(time.Second / 10), + Resolver: net.LookupAddr, } go r.loop() return &r @@ -43,43 +40,37 @@ func newReverseResolver() *reverseResolver { // Get the reverse resolution for an IP address if already in the cache, // a gcache.NotFoundKeyError error otherwise. // Note: it returns one of the possible names that can be obtained for that IP. -func (r *reverseResolver) Get(address string, wait bool) (string, error) { +func (r *ReverseResolver) Get(address string) (string, error) { val, err := r.cache.Get(address) if err == nil { return val.(string), nil } if err == gcache.NotFoundKeyError { - request := revResRequest{address: address, done: make(chan struct{})} // we trigger a asynchronous reverse resolution when not cached select { - case r.addresses <- request: - if wait { - <-request.done - } + case r.addresses <- address: default: } } return "", err } -func (r *reverseResolver) loop() { - throttle := time.Tick(time.Second / 10) +func (r *ReverseResolver) loop() { for request := range r.addresses { - <-throttle // rate limit our DNS resolutions + <-r.Throttle // rate limit our DNS resolutions // and check if the answer is already in the cache - if _, err := r.cache.Get(request.address); err == nil { + if _, err := r.cache.Get(request); err == nil { continue } - names, err := r.resolver(request.address) + names, err := r.Resolver(request) if err == nil && len(names) > 0 { name := strings.TrimRight(names[0], ".") - r.cache.Set(request.address, name) + r.cache.Set(request, name) } - close(request.done) } } // Stop the async reverse resolver -func (r *reverseResolver) Stop() { +func (r *ReverseResolver) Stop() { close(r.addresses) } diff --git a/probe/endpoint/resolver_internal_test.go b/probe/endpoint/resolver_internal_test.go deleted file mode 100644 index e9a2d264e..000000000 --- a/probe/endpoint/resolver_internal_test.go +++ /dev/null @@ -1,41 +0,0 @@ -package endpoint - -import ( - "errors" - "testing" -) - -func TestReverseResolver(t *testing.T) { - tests := map[string]string{ - "8.8.8.8": "google-public-dns-a.google.com", - "8.8.4.4": "google-public-dns-b.google.com", - } - - revRes := newReverseResolver() - - // use a mocked resolver function - revRes.resolver = func(addr string) (names []string, err error) { - if name, ok := tests[addr]; ok { - return []string{name}, nil - } - return []string{}, errors.New("invalid IP") - } - - // first time: no names are returned for our reverse resolutions - for ip := range tests { - if have, err := revRes.Get(ip, true); have != "" || err == nil { - t.Errorf("we didn't get an error, or the cache was not empty, when trying to resolve '%q'", ip) - } - } - - // so, if we check again these IPs, we should have the names now - for ip, want := range tests { - have, err := revRes.Get(ip, true) - if err != nil { - t.Errorf("%s: %v", ip, err) - } - if want != have { - t.Errorf("%s: want %q, have %q", ip, want, have) - } - } -} diff --git a/probe/endpoint/resolver_test.go b/probe/endpoint/resolver_test.go new file mode 100644 index 000000000..31c061e7f --- /dev/null +++ b/probe/endpoint/resolver_test.go @@ -0,0 +1,38 @@ +package endpoint_test + +import ( + "errors" + "testing" + "time" + + . "github.com/weaveworks/scope/probe/endpoint" + "github.com/weaveworks/scope/test" +) + +func TestReverseResolver(t *testing.T) { + tests := map[string]string{ + "1.2.3.4": "test.domain.name", + "4.3.2.1": "im.a.little.tea.pot", + } + + revRes := NewReverseResolver() + defer revRes.Stop() + + // use a mocked resolver function + revRes.Resolver = func(addr string) (names []string, err error) { + if name, ok := tests[addr]; ok { + return []string{name}, nil + } + return []string{}, errors.New("invalid IP") + } + + // Up the rate limit so the test runs faster + revRes.Throttle = time.Tick(time.Millisecond) + + for ip, hostname := range tests { + test.Poll(t, 100*time.Millisecond, hostname, func() interface{} { + result, _ := revRes.Get(ip) + return result + }) + } +} From 7513b4e39654b68932f771a607aeab8636d2c8f1 Mon Sep 17 00:00:00 2001 From: Peter Bourgon Date: Mon, 7 Sep 2015 10:36:35 +0200 Subject: [PATCH 3/3] Wrap comments at 80col throughout the fileset --- probe/endpoint/reporter.go | 12 ++++++--- probe/endpoint/resolver.go | 16 ++++++------ probe/endpoint/resolver_test.go | 4 +-- report/topology.go | 43 ++++++++++++++++++--------------- 4 files changed, 41 insertions(+), 34 deletions(-) diff --git a/probe/endpoint/reporter.go b/probe/endpoint/reporter.go index 8685fa932..04e6dc71b 100644 --- a/probe/endpoint/reporter.go +++ b/probe/endpoint/reporter.go @@ -148,7 +148,8 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin }) ) - // in case we have a reverse resolution for the IP, we can use it for the name... + // In case we have a reverse resolution for the IP, we can use it for + // the name... if revRemoteName, err := r.revResolver.Get(remoteAddr); err == nil { remoteNode = remoteNode.AddMetadata(map[string]string{ "name": revRemoteName, @@ -156,7 +157,8 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin } if localIsClient { - // New nodes are merged into the report so we don't need to do any counting here; the merge does it for us. + // New nodes are merged into the report so we don't need to do any + // counting here; the merge does it for us. localNode = localNode.WithEdge(remoteAddressNodeID, report.EdgeMetadata{ MaxConnCountTCP: newu64(1), }) @@ -187,7 +189,8 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin }) ) - // in case we have a reverse resolution for the IP, we can use it for the name... + // In case we have a reverse resolution for the IP, we can use it for + // the name... if revRemoteName, err := r.revResolver.Get(remoteAddr); err == nil { remoteNode = remoteNode.AddMetadata(map[string]string{ "name": revRemoteName, @@ -195,7 +198,8 @@ func (r *Reporter) addConnection(rpt *report.Report, localAddr, remoteAddr strin } if localIsClient { - // New nodes are merged into the report so we don't need to do any counting here; the merge does it for us. + // New nodes are merged into the report so we don't need to do any + // counting here; the merge does it for us. localNode = localNode.WithEdge(remoteEndpointNodeID, report.EdgeMetadata{ MaxConnCountTCP: newu64(1), }) diff --git a/probe/endpoint/resolver.go b/probe/endpoint/resolver.go index 210e93832..9f2609adc 100644 --- a/probe/endpoint/resolver.go +++ b/probe/endpoint/resolver.go @@ -16,7 +16,7 @@ const ( type revResFunc func(addr string) (names []string, err error) -// ReverseResolver is a caching, reverse resolver +// ReverseResolver is a caching, reverse resolver. type ReverseResolver struct { addresses chan string cache gcache.Cache @@ -24,8 +24,8 @@ type ReverseResolver struct { Resolver revResFunc } -// NewReverseResolver starts a new reverse resolver that -// performs reverse resolutions and caches the result. +// NewReverseResolver starts a new reverse resolver that performs reverse +// resolutions and caches the result. func NewReverseResolver() *ReverseResolver { r := ReverseResolver{ addresses: make(chan string, rAddrBacklog), @@ -37,16 +37,16 @@ func NewReverseResolver() *ReverseResolver { return &r } -// Get the reverse resolution for an IP address if already in the cache, -// a gcache.NotFoundKeyError error otherwise. -// Note: it returns one of the possible names that can be obtained for that IP. +// Get the reverse resolution for an IP address if already in the cache, a +// gcache.NotFoundKeyError error otherwise. Note: it returns one of the +// possible names that can be obtained for that IP. func (r *ReverseResolver) Get(address string) (string, error) { val, err := r.cache.Get(address) if err == nil { return val.(string), nil } if err == gcache.NotFoundKeyError { - // we trigger a asynchronous reverse resolution when not cached + // We trigger a asynchronous reverse resolution when not cached select { case r.addresses <- address: default: @@ -70,7 +70,7 @@ func (r *ReverseResolver) loop() { } } -// Stop the async reverse resolver +// Stop the async reverse resolver. func (r *ReverseResolver) Stop() { close(r.addresses) } diff --git a/probe/endpoint/resolver_test.go b/probe/endpoint/resolver_test.go index 31c061e7f..25a43cc3c 100644 --- a/probe/endpoint/resolver_test.go +++ b/probe/endpoint/resolver_test.go @@ -18,7 +18,7 @@ func TestReverseResolver(t *testing.T) { revRes := NewReverseResolver() defer revRes.Stop() - // use a mocked resolver function + // Use a mocked resolver function. revRes.Resolver = func(addr string) (names []string, err error) { if name, ok := tests[addr]; ok { return []string{name}, nil @@ -26,7 +26,7 @@ func TestReverseResolver(t *testing.T) { return []string{}, errors.New("invalid IP") } - // Up the rate limit so the test runs faster + // Up the rate limit so the test runs faster. revRes.Throttle = time.Tick(time.Millisecond) for ip, hostname := range tests { diff --git a/report/topology.go b/report/topology.go index bf237d996..263cc8c1f 100644 --- a/report/topology.go +++ b/report/topology.go @@ -6,9 +6,9 @@ import ( ) // Topology describes a specific view of a network. It consists of nodes and -// edges, and metadata about those nodes and edges, represented by EdgeMetadatas -// and Nodes respectively. Edges are directional, and embedded in the -// Node struct. +// edges, and metadata about those nodes and edges, represented by +// EdgeMetadatas and Nodes respectively. Edges are directional, and embedded +// in the Node struct. type Topology struct { Nodes } @@ -20,8 +20,9 @@ func MakeTopology() Topology { } } -// WithNode produces a topology from t, with nmd added under key nodeID; if a node already exists -// for this key, nmd is merged with that node. NB A fresh topology is returned. +// WithNode produces a topology from t, with nmd added under key nodeID; if a +// node already exists for this key, nmd is merged with that node. Note that a +// fresh topology is returned. func (t Topology) WithNode(nodeID string, nmd Node) Topology { if existing, ok := t.Nodes[nodeID]; ok { nmd = nmd.Merge(existing) @@ -70,9 +71,9 @@ func (n Nodes) Merge(other Nodes) Nodes { return cp } -// Node describes a superset of the metadata that probes can collect -// about a given node in a given topology, along with the edges emanating -// from the node and metadata about those edges. +// Node describes a superset of the metadata that probes can collect about a +// given node in a given topology, along with the edges emanating from the +// node and metadata about those edges. type Node struct { Metadata `json:"-"` Counters `json:"-"` @@ -102,20 +103,21 @@ func (n Node) WithMetadata(m map[string]string) Node { return result } -// AddMetadata returns a fresh copy of n, with Metadata set to the merge of n and the metadata provided +// AddMetadata returns a fresh copy of n, with Metadata set to the merge of n +// and the metadata provided. func (n Node) AddMetadata(m map[string]string) Node { additional := MakeNodeWith(m) return n.Merge(additional) } -// WithCounters returns a fresh copy of n, with Counters set to c +// WithCounters returns a fresh copy of n, with Counters set to c. func (n Node) WithCounters(c map[string]int) Node { result := n.Copy() result.Counters = c return result } -// WithAdjacency returns a fresh copy of n, with Adjacency set to a +// WithAdjacency returns a fresh copy of n, with Adjacency set to a. func (n Node) WithAdjacency(a IDList) Node { result := n.Copy() result.Adjacency = a @@ -129,7 +131,8 @@ func (n Node) WithAdjacent(a string) Node { return result } -// WithEdge returns a fresh copy of n, with 'dst' added to Adjacency and md added to EdgeMetadata +// WithEdge returns a fresh copy of n, with 'dst' added to Adjacency and md +// added to EdgeMetadata. func (n Node) WithEdge(dst string, md EdgeMetadata) Node { result := n.Copy() result.Adjacency = result.Adjacency.Add(dst) @@ -158,7 +161,7 @@ func (n Node) Merge(other Node) Node { return cp } -// Metadata is a string->string map +// Metadata is a string->string map. type Metadata map[string]string // Merge merges two node metadata maps together. In case of conflict, the @@ -172,7 +175,7 @@ func (m Metadata) Merge(other Metadata) Metadata { return result } -// Copy creates a deep copy of the Metadata +// Copy creates a deep copy of the Metadata. func (m Metadata) Copy() Metadata { result := Metadata{} for k, v := range m { @@ -181,11 +184,11 @@ func (m Metadata) Copy() Metadata { return result } -// Counters is a string->int map +// Counters is a string->int map. type Counters map[string]int -// Merge merges two sets of counters into a fresh set of counters, -// summing values where appropriate +// Merge merges two sets of counters into a fresh set of counters, summing +// values where appropriate. func (c Counters) Merge(other Counters) Counters { result := c.Copy() for k, v := range other { @@ -194,7 +197,7 @@ func (c Counters) Merge(other Counters) Counters { return result } -// Copy creates a deep copy of the Counters +// Copy creates a deep copy of the Counters. func (c Counters) Copy() Counters { result := Counters{} for k, v := range c { @@ -203,8 +206,8 @@ func (c Counters) Copy() Counters { return result } -// EdgeMetadatas collect metadata about each edge in a topology. Keys are -// the remote node IDs, as in Adjacency. +// EdgeMetadatas collect metadata about each edge in a topology. Keys are the +// remote node IDs, as in Adjacency. type EdgeMetadatas map[string]EdgeMetadata // Copy returns a value copy of the EdgeMetadatas.