Only fetch weave status report once per tick.

This commit is contained in:
Tom Wilkie
2015-09-08 18:16:03 +00:00
parent 5bd324db3f
commit b7c22b7a8f
6 changed files with 35 additions and 24 deletions

View File

@@ -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))

View File

@@ -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,

View File

@@ -30,6 +30,8 @@ func TestWeaveTaggerOverlayTopology(t *testing.T) {
t.Fatal(err)
}
w.Tick()
{
have, err := w.Report()
if err != nil {

View File

@@ -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)

View File

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

View File

@@ -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