From e3c5b7f36d98308e58ac2a5d45c73a3dd9fde065 Mon Sep 17 00:00:00 2001 From: Peter Bourgon Date: Thu, 11 Jun 2015 18:45:28 +0200 Subject: [PATCH] Add WeaveTagger - report: add Overlay topology - probe/tag: introduce WeaveTagger --- probe/main.go | 20 ++++++- probe/tag/weave_tagger.go | 102 +++++++++++++++++++++++++++++++++ probe/tag/weave_tagger_test.go | 51 +++++++++++++++++ report/id.go | 6 ++ report/merge.go | 1 + report/report.go | 7 +++ 6 files changed, 185 insertions(+), 2 deletions(-) create mode 100644 probe/tag/weave_tagger.go create mode 100644 probe/tag/weave_tagger_test.go diff --git a/probe/main.go b/probe/main.go index 9b04b9763..544ae28d1 100644 --- a/probe/main.go +++ b/probe/main.go @@ -33,6 +33,7 @@ func main() { spyProcs = flag.Bool("processes", true, "report processes (needs root)") dockerEnabled = flag.Bool("docker", true, "collect Docker-related attributes for processes") dockerInterval = flag.Duration("docker.interval", 10*time.Second, "how often to update Docker attributes") + weaveRouterAddr = flag.String("weave.router.addr", "", "IP address or FQDN of the Weave router") procRoot = flag.String("proc.root", "/proc", "location of the proc filesystem") ) flag.Parse() @@ -71,9 +72,11 @@ func main() { hostID = hostName // TODO: we should sanitize the hostname ) + var ( + dockerTagger *tag.DockerTagger + weaveTagger *tag.WeaveTagger + ) taggers := []tag.Tagger{tag.NewTopologyTagger(), tag.NewOriginHostTagger(hostID)} - - var dockerTagger *tag.DockerTagger if *dockerEnabled && runtime.GOOS == linux { var err error dockerTagger, err = tag.NewDockerTagger(*procRoot, *dockerInterval) @@ -84,6 +87,15 @@ func main() { taggers = append(taggers, dockerTagger) } + if *weaveRouterAddr != "" { + var err error + weaveTagger, err = tag.NewWeaveTagger(*weaveRouterAddr) + if err != nil { + log.Fatalf("failed to start Weave tagger: %v", err) + } + taggers = append(taggers, weaveTagger) + } + log.Printf("listening on %s", *listen) quit := make(chan struct{}) @@ -120,6 +132,10 @@ func main() { r.Container.Merge(dockerTagger.ContainerTopology(hostID)) } + if weaveTagger != nil { + r.Overlay.Merge(weaveTagger.OverlayTopology()) + } + r = tag.Apply(r, taggers) case <-quit: diff --git a/probe/tag/weave_tagger.go b/probe/tag/weave_tagger.go new file mode 100644 index 000000000..1865bc37e --- /dev/null +++ b/probe/tag/weave_tagger.go @@ -0,0 +1,102 @@ +package tag + +import ( + "encoding/json" + "fmt" + "log" + "net" + "net/http" + "net/url" + "strings" + + "github.com/weaveworks/scope/report" +) + +const ( + // WeavePeerName is the key for the peer name, typically a MAC address. + // It's also the unscoped + WeavePeerName = "weave_peer_name" + + // WeavePeerNickName is the key for the peer nickname, typically a + // hostname. + WeavePeerNickName = "weave_peer_nick_name" +) + +// WeaveTagger represents a single Weave router, presumably on the same host +// as the probe. It can produce an Overlay topology, and in theory can tag +// existing topologies with foreign keys to overlay -- though I'm not sure +// what that would look like in practice right now. +type WeaveTagger struct { + url string +} + +// NewWeaveTagger returns a new Weave tagger based on the Weave router at +// address. The address should be an IP or FQDN, no port. +func NewWeaveTagger(weaveRouterAddress string) (*WeaveTagger, error) { + s, err := sanitize("http://", 6784, "/status-json")(weaveRouterAddress) + if err != nil { + return nil, err + } + return &WeaveTagger{s}, nil +} + +// Tag implements Tagger. +func (t WeaveTagger) Tag(r report.Report) report.Report { + // The status-json endpoint doesn't return any link information, so + // there's nothing to tag, yet. + return r +} + +// OverlayTopology produces an overlay topology from the Weave router. +func (t WeaveTagger) OverlayTopology() report.Topology { + topology := report.NewTopology() + + resp, err := http.Get(t.url) + if err != nil { + log.Printf("Weave Tagger: %v", err) + return topology + } + defer resp.Body.Close() + + var status struct { + Peers []struct { + Name string `json:"Name"` + NickName string `json:"NickName"` + } `json:"Peers"` + } + if err := json.NewDecoder(resp.Body).Decode(&status); err != nil { + log.Printf("Weave Tagger: %v", err) + return topology + } + + for _, peer := range status.Peers { + topology.NodeMetadatas[report.MakeOverlayNodeID(peer.Name)] = report.NodeMetadata{ + WeavePeerName: peer.Name, + WeavePeerNickName: peer.NickName, + } + } + + return topology +} + +func sanitize(scheme string, port int, path string) func(string) (string, error) { + return func(s string) (string, error) { + if s == "" { + return "", fmt.Errorf("no host") + } + if !strings.HasPrefix(s, "http") { + s = scheme + s + } + u, err := url.Parse(s) + if err != nil { + return "", err + } + if _, _, err = net.SplitHostPort(u.Host); err != nil { + u.Host += fmt.Sprintf(":%d", port) + } + if u.Path != path { + u.Path = path + } + return u.String(), nil + } +} diff --git a/probe/tag/weave_tagger_test.go b/probe/tag/weave_tagger_test.go new file mode 100644 index 000000000..fdfc251cf --- /dev/null +++ b/probe/tag/weave_tagger_test.go @@ -0,0 +1,51 @@ +package tag_test + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "reflect" + "testing" + + "github.com/weaveworks/scope/probe/tag" + "github.com/weaveworks/scope/report" +) + +func TestWeaveTaggerOverlayTopology(t *testing.T) { + s := httptest.NewServer(http.HandlerFunc(mockWeaveRouter)) + defer s.Close() + + w, err := tag.NewWeaveTagger(s.URL) + if err != nil { + t.Fatal(err) + } + + if want, have := (report.Topology{ + Adjacency: report.Adjacency{}, + EdgeMetadatas: report.EdgeMetadatas{}, + NodeMetadatas: report.NodeMetadatas{ + report.MakeOverlayNodeID(mockWeavePeerName): { + tag.WeavePeerName: mockWeavePeerName, + tag.WeavePeerNickName: mockWeavePeerNickName, + }, + }, + }), w.OverlayTopology(); !reflect.DeepEqual(want, have) { + t.Errorf("want\n\t%#v, have\n\t%#v", want, have) + } +} + +const ( + mockWeavePeerName = "winnebago" + mockWeavePeerNickName = "winny" +) + +func mockWeaveRouter(w http.ResponseWriter, r *http.Request) { + if err := json.NewEncoder(w).Encode(map[string]interface{}{ + "Peers": []map[string]interface{}{{ + "Name": mockWeavePeerName, + "NickName": mockWeavePeerNickName, + }}, + }); err != nil { + println(err.Error()) + } +} diff --git a/report/id.go b/report/id.go index 5973a51c1..df7c300e1 100644 --- a/report/id.go +++ b/report/id.go @@ -81,6 +81,12 @@ func MakeContainerNodeID(hostID, containerID string) string { return hostID + ScopeDelim + containerID } +// MakeOverlayNodeID produces an overlay topology node ID from a router peer's +// name, which is assumed to be globally unique. +func MakeOverlayNodeID(peerName string) string { + return "#" + peerName +} + // ParseNodeID produces the host ID and remainder (typically an address) from // a node ID. Note that hostID may be blank. func ParseNodeID(nodeID string) (hostID string, remainder string, ok bool) { diff --git a/report/merge.go b/report/merge.go index a4e95401e..87f85a3ce 100644 --- a/report/merge.go +++ b/report/merge.go @@ -11,6 +11,7 @@ func (r *Report) Merge(other Report) { r.Process.Merge(other.Process) r.Container.Merge(other.Container) r.Host.Merge(other.Host) + r.Overlay.Merge(other.Overlay) } // Merge merges another Topology into the receiver. diff --git a/report/report.go b/report/report.go index 3a6fdad1c..36dfdb51f 100644 --- a/report/report.go +++ b/report/report.go @@ -31,6 +31,11 @@ type Report struct { // like operating system, load, etc. The information is scraped by the // probes with each published report. Edges are not present. Host Topology + + // Overlay nodes are active peers in any software-defined network that's + // overlaid on the infrastructure. The information is scraped by polling + // their status endpoints. Edges could be present, but aren't currently. + Overlay Topology } const ( @@ -90,6 +95,7 @@ func MakeReport() Report { Process: NewTopology(), Container: NewTopology(), Host: NewTopology(), + Overlay: NewTopology(), } } @@ -102,6 +108,7 @@ func (r Report) Squash() Report { r.Process = r.Process.Squash(PanicIDAddresser, localNetworks) r.Container = r.Container.Squash(PanicIDAddresser, localNetworks) r.Host = r.Host.Squash(PanicIDAddresser, localNetworks) + r.Overlay = r.Overlay.Squash(PanicIDAddresser, localNetworks) return r }