Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Refactoring of probe main #633

Merged
merged 2 commits into from
Nov 9, 2015
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ app/app
app/scope-app
probe/probe
probe/scope-probe
prog/probe/scope-probe
docker/scope-app
docker/scope-probe
docker/docker*
Expand Down
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
SUDO=sudo -E
DOCKERHUB_USER=weaveworks
APP_EXE=app/scope-app
PROBE_EXE=probe/scope-probe
PROBE_EXE=prog/probe/scope-probe
FIXPROBE_EXE=experimental/fixprobe/fixprobe
SCOPE_IMAGE=$(DOCKERHUB_USER)/scope
SCOPE_EXPORT=scope.tar
Expand Down
4 changes: 2 additions & 2 deletions circle.yml
Original file line number Diff line number Diff line change
Expand Up @@ -45,9 +45,9 @@ test:
parallel: true
- cd $SRCDIR; make RM= static:
parallel: true
- cd $SRCDIR; rm -f app/scope-app probe/scope-probe; if [ "$CIRCLE_NODE_INDEX" = "0" ]; then GOARCH=arm make RM= app/scope-app probe/scope-probe; else GOOS=darwin make RM= app/scope-app probe/scope-probe; fi:
- cd $SRCDIR; rm -f app/scope-app prog/probe/scope-probe; if [ "$CIRCLE_NODE_INDEX" = "0" ]; then GOARCH=arm make RM= app/scope-app prog/probe/scope-probe; else GOOS=darwin make RM= app/scope-app prog/probe/scope-probe; fi:
parallel: true
- cd $SRCDIR; rm -f app/scope-app probe/scope-probe; make RM=:
- cd $SRCDIR; rm -f app/scope-app prog/probe/scope-probe; make RM=:
parallel: true
- cd $SRCDIR/experimental; ./build_on_circle.sh:
parallel: true
Expand Down
5 changes: 3 additions & 2 deletions probe/hostname.go
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
package main
package probe

import "os"

func hostname() string {
// Hostname returns the hostname of this host.
func Hostname() string {
if hostname := os.Getenv("SCOPE_HOSTNAME"); hostname != "" {
return hostname
}
Expand Down
27 changes: 0 additions & 27 deletions probe/instrumentation.go

This file was deleted.

162 changes: 162 additions & 0 deletions probe/probe.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
package probe

import (
"log"
"sync"
"time"

"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/xfer"
)

// Probe sits there, generating and publishing reports.
type Probe struct {
spyInterval, publishInterval time.Duration
publisher xfer.Publisher

tickers []Ticker
reporters []Reporter
taggers []Tagger

quit chan struct{}
done sync.WaitGroup
rpt syncReport
}

// Tagger tags nodes with value-add node metadata.
type Tagger interface {
Tag(r report.Report) (report.Report, error)
}

// Reporter generates Reports.
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
}

// New makes a new Probe.
func New(spyInterval, publishInterval time.Duration, publisher xfer.Publisher) *Probe {
result := &Probe{
spyInterval: spyInterval,
publishInterval: publishInterval,
publisher: publisher,
quit: make(chan struct{}),
}
result.rpt.swap(report.MakeReport())
return result
}

// AddTagger adds a new Tagger to the Probe
func (p *Probe) AddTagger(ts ...Tagger) {
p.taggers = append(p.taggers, ts...)
}

// AddReporter adds a new Reported to the Probe
func (p *Probe) AddReporter(rs ...Reporter) {
p.reporters = append(p.reporters, rs...)
}

// AddTicker adds a new Ticker to the Probe
func (p *Probe) AddTicker(ts ...Ticker) {
p.tickers = append(p.tickers, ts...)
}

// Start starts the probe
func (p *Probe) Start() {
p.done.Add(2)
go p.spyLoop()
go p.publishLoop()
}

// Stop stops the probe
func (p *Probe) Stop() {
close(p.quit)
p.done.Wait()
}

func (p *Probe) spyLoop() {
defer p.done.Done()
spyTick := time.Tick(p.spyInterval)

for {
select {
case <-spyTick:
start := time.Now()
for _, ticker := range p.tickers {
if err := ticker.Tick(); err != nil {
log.Printf("error doing ticker: %v", err)
}
}

localReport := p.rpt.copy()
localReport = localReport.Merge(p.report())
localReport = p.tag(localReport)
p.rpt.swap(localReport)

if took := time.Since(start); took > p.spyInterval {
log.Printf("report generation took too long (%s)", took)
}

case <-p.quit:
return
}
}
}

func (p *Probe) report() report.Report {
reports := make(chan report.Report, len(p.reporters))
for _, rep := range p.reporters {
go func(rep Reporter) {
newReport, err := rep.Report()
if err != nil {
log.Printf("error generating report: %v", err)
newReport = report.MakeReport() // empty is OK to merge
}
reports <- newReport
}(rep)
}

result := report.MakeReport()
for i := 0; i < cap(reports); i++ {
result = result.Merge(<-reports)
}
return result
}

func (p *Probe) tag(r report.Report) report.Report {
var err error
for _, tagger := range p.taggers {
r, err = tagger.Tag(r)
if err != nil {
log.Printf("error applying tagger: %v", err)
}
}
return r
}

func (p *Probe) publishLoop() {
defer p.done.Done()
var (
pubTick = time.Tick(p.publishInterval)
rptPub = xfer.NewReportPublisher(p.publisher)
)

for {
select {
case <-pubTick:
localReport := p.rpt.swap(report.MakeReport())
if err := rptPub.Publish(localReport); err != nil {
log.Printf("publish: %v", err)
}

case <-p.quit:
return
}
}
}
86 changes: 86 additions & 0 deletions probe/probe_internal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
package probe

import (
"compress/gzip"
"encoding/gob"
"io"
"reflect"
"testing"
"time"

"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/test"
)

func TestApply(t *testing.T) {
var (
endpointNodeID = "c"
addressNodeID = "d"
endpointNode = report.MakeNodeWith(map[string]string{"5": "6"})
addressNode = report.MakeNodeWith(map[string]string{"7": "8"})
)

p := New(0, 0, nil)
p.AddTagger(NewTopologyTagger())

r := report.MakeReport()
r.Endpoint.AddNode(endpointNodeID, endpointNode)
r.Address.AddNode(addressNodeID, addressNode)
r = p.tag(r)

for _, tuple := range []struct {
want report.Node
from report.Topology
via string
}{
{endpointNode.Merge(report.MakeNodeWith(map[string]string{"topology": "endpoint"})), r.Endpoint, endpointNodeID},
{addressNode.Merge(report.MakeNodeWith(map[string]string{"topology": "address"})), r.Address, addressNodeID},
} {
if want, have := tuple.want, tuple.from.Nodes[tuple.via]; !reflect.DeepEqual(want, have) {
t.Errorf("want %+v, have %+v", want, have)
}
}
}

type mockReporter struct {
r report.Report
}

func (m mockReporter) Report() (report.Report, error) {
return m.r.Copy(), nil
}

type mockPublisher struct {
have chan report.Report
}

func (m mockPublisher) Publish(in io.Reader) error {
var r report.Report
if reader, err := gzip.NewReader(in); err != nil {
return err
} else if err := gob.NewDecoder(reader).Decode(&r); err != nil {
return err
}
m.have <- r
return nil
}

func (m mockPublisher) Stop() {
close(m.have)
}

func TestProbe(t *testing.T) {
want := report.MakeReport()
node := report.MakeNodeWith(map[string]string{"b": "c"})
want.Endpoint.AddNode("a", node)
pub := mockPublisher{make(chan report.Report)}

p := New(10*time.Millisecond, 100*time.Millisecond, pub)
p.AddReporter(mockReporter{want})
p.Start()
defer p.Stop()

test.Poll(t, 300*time.Millisecond, want, func() interface{} {
return <-pub.have
})
}
26 changes: 26 additions & 0 deletions probe/sync_report.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package probe

import (
"sync"

"github.com/weaveworks/scope/report"
)

type syncReport struct {
mtx sync.RWMutex
rpt report.Report
}

func (r *syncReport) swap(other report.Report) report.Report {
r.mtx.Lock()
defer r.mtx.Unlock()
old := r.rpt
r.rpt = other
return old
}

func (r *syncReport) copy() report.Report {
r.mtx.RLock()
defer r.mtx.RUnlock()
return r.rpt.Copy()
}
45 changes: 0 additions & 45 deletions probe/tag_report_test.go

This file was deleted.

Loading