diff --git a/probe/main.go b/probe/main.go index 385f82f07..161276719 100644 --- a/probe/main.go +++ b/probe/main.go @@ -113,6 +113,7 @@ func main() { var ( endpointReporter = endpoint.NewReporter(hostID, hostName, *spyProcs, *useConntrack) processCache = process.NewCachingWalker(process.NewWalker(*procRoot)) + tickers = []Ticker{processCache} reporters = []Reporter{ endpointReporter, host.NewReporter(hostID, hostName, localNets), @@ -142,6 +143,7 @@ func main() { if err != nil { log.Fatalf("failed to start Weave tagger: %v", err) } + tickers = append(tickers, weave) taggers = append(taggers, weave) reporters = append(reporters, weave) } @@ -187,8 +189,11 @@ func main() { case <-spyTick: start := time.Now() - if err := processCache.Update(); err != nil { - log.Printf("error reading processes: %v", err) + + for _, ticker := range tickers { + if err := ticker.Tick(); err != nil { + log.Printf("error doing ticker: %v", err) + } } r = r.Merge(doReport(reporters)) diff --git a/probe/overlay/weave.go b/probe/overlay/weave.go index b9699cdd3..42d3f55e5 100644 --- a/probe/overlay/weave.go +++ b/probe/overlay/weave.go @@ -43,6 +43,7 @@ var ipMatch = regexp.MustCompile(`([0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3 type Weave struct { url string hostID string + status weaveStatus } type weaveStatus struct { @@ -75,24 +76,30 @@ func NewWeave(hostID, weaveRouterAddress string) (*Weave, error) { }, nil } -func (w Weave) update() (weaveStatus, error) { +// Tick implements Ticker +func (w *Weave) Tick() error { var result weaveStatus req, err := http.NewRequest("GET", w.url, nil) if err != nil { - return result, err + return err } req.Header.Add("Accept", "application/json") resp, err := http.DefaultClient.Do(req) if err != nil { - return result, err + return err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return result, fmt.Errorf("Weave Tagger: got %d", resp.StatusCode) + return fmt.Errorf("Weave Tagger: got %d", resp.StatusCode) } - return result, json.NewDecoder(resp.Body).Decode(&result) + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return err + } + + w.status = result + return nil } type psEntry struct { @@ -132,7 +139,7 @@ func (w Weave) ps() ([]psEntry, error) { return result, scanner.Err() } -func (w Weave) tagContainer(r report.Report, containerIDPrefix, macAddress string, ips []string) { +func (w *Weave) tagContainer(r report.Report, containerIDPrefix, macAddress string, ips []string) { for nodeid, nmd := range r.Container.Nodes { idPrefix := nmd.Metadata[docker.ContainerID][:12] if idPrefix != containerIDPrefix { @@ -150,12 +157,7 @@ func (w Weave) tagContainer(r report.Report, containerIDPrefix, macAddress strin // Tag implements Tagger. func (w Weave) Tag(r report.Report) (report.Report, error) { - status, err := w.update() - if err != nil { - return r, nil - } - - for _, entry := range status.DNS.Entries { + for _, entry := range w.status.DNS.Entries { if entry.Tombstone > 0 { continue } @@ -183,12 +185,7 @@ func (w Weave) Tag(r report.Report) (report.Report, error) { // Report implements Reporter. func (w Weave) Report() (report.Report, error) { r := report.MakeReport() - status, err := w.update() - if err != nil { - return r, err - } - - for _, peer := range status.Router.Peers { + for _, peer := range w.status.Router.Peers { r.Overlay.Nodes[report.MakeOverlayNodeID(peer.Name)] = report.MakeNodeWith(map[string]string{ WeavePeerName: peer.Name, WeavePeerNickName: peer.NickName, diff --git a/probe/overlay/weave_test.go b/probe/overlay/weave_test.go index f5d302a65..bb5f18f6e 100644 --- a/probe/overlay/weave_test.go +++ b/probe/overlay/weave_test.go @@ -30,6 +30,8 @@ func TestWeaveTaggerOverlayTopology(t *testing.T) { t.Fatal(err) } + w.Tick() + { have, err := w.Report() if err != nil { diff --git a/probe/process/walker.go b/probe/process/walker.go index 4a88c45cd..4ba20027e 100644 --- a/probe/process/walker.go +++ b/probe/process/walker.go @@ -39,8 +39,8 @@ func (c *CachingWalker) Walk(f func(Process)) error { return nil } -// Update updates cached copy of process list -func (c *CachingWalker) Update() error { +// Tick updates cached copy of process list +func (c *CachingWalker) Tick() error { newCache := []Process{} err := c.source.Walk(func(p Process) { newCache = append(newCache, p) diff --git a/probe/process/walker_test.go b/probe/process/walker_test.go index 37d445891..6046006bb 100644 --- a/probe/process/walker_test.go +++ b/probe/process/walker_test.go @@ -29,7 +29,7 @@ func TestCache(t *testing.T) { processes: processes, } cachingWalker := process.NewCachingWalker(walker) - err := cachingWalker.Update() + err := cachingWalker.Tick() if err != nil { t.Fatal(err) } @@ -45,7 +45,7 @@ func TestCache(t *testing.T) { t.Errorf("%v (%v)", test.Diff(processes, have), err) } - err = cachingWalker.Update() + err = cachingWalker.Tick() if err != nil { t.Fatal(err) } diff --git a/probe/tag_report.go b/probe/tag_report.go index 15024d77a..5a3bfb4e5 100644 --- a/probe/tag_report.go +++ b/probe/tag_report.go @@ -16,6 +16,13 @@ type Reporter interface { Report() (report.Report, error) } +// Ticker is something which will be invoked every spyDuration. +// It's useful for things that should be updated on that interval. +// For example, cached shared state between Taggers and Reporters. +type Ticker interface { + Tick() error +} + // Apply tags the report with all the taggers. func Apply(r report.Report, taggers []Tagger) report.Report { var err error