From 92b24793f1ba1c3f8bfbff9f6f39722001bd7a1f Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Mon, 9 Nov 2015 16:14:22 +0000 Subject: [PATCH 1/3] Push mini-reports when container's state changes. --- probe/docker/registry.go | 18 ++++++++++++++ probe/docker/reporter.go | 20 +++++++++++++++- probe/docker/reporter_test.go | 4 +++- probe/probe.go | 45 ++++++++++++++++++++++++++++------- probe/sync_report.go | 26 -------------------- prog/probe/main.go | 5 +++- 6 files changed, 81 insertions(+), 37 deletions(-) delete mode 100644 probe/sync_report.go diff --git a/probe/docker/registry.go b/probe/docker/registry.go index 3c21bd4e7..e8ad506ef 100644 --- a/probe/docker/registry.go +++ b/probe/docker/registry.go @@ -31,14 +31,19 @@ type Registry interface { LockedPIDLookup(f func(func(int) Container)) WalkContainers(f func(Container)) WalkImages(f func(*docker_client.APIImages)) + WatchContainerUpdates(ContainerUpdateWatcher) } +// ContainerUpdateWatcher is the type of functions that get called when containers are updated. +type ContainerUpdateWatcher func(c Container) + type registry struct { sync.RWMutex quit chan chan struct{} interval time.Duration client Client + watchers []ContainerUpdateWatcher containers map[string]Container containersByPID map[int]Container images map[string]*docker_client.APIImages @@ -91,6 +96,14 @@ func (r *registry) Stop() { <-ch } +// WatchContainerUpdates registers a callback to be called +// whenever a container is updated. +func (r *registry) WatchContainerUpdates(f ContainerUpdateWatcher) { + r.Lock() + defer r.Unlock() + r.watchers = append(r.watchers, f) +} + func (r *registry) loop() { for { // NB listenForEvents blocks. @@ -244,6 +257,11 @@ func (r *registry) updateContainerState(containerID string) { c.UpdateState(dockerContainer) } + // Trigger anyone watching for updates + for _, f := range r.watchers { + f(c) + } + // And finally, ensure we gather stats for it if dockerContainer.State.Running { if err := c.StartGatheringStats(); err != nil { diff --git a/probe/docker/reporter.go b/probe/docker/reporter.go index 8f7561efc..a2c10c19e 100644 --- a/probe/docker/reporter.go +++ b/probe/docker/reporter.go @@ -1,10 +1,12 @@ package docker import ( + "log" "net" docker_client "github.com/fsouza/go-dockerclient" + "github.com/weaveworks/scope/probe" "github.com/weaveworks/scope/report" ) @@ -18,16 +20,32 @@ const ( type Reporter struct { registry Registry hostID string + probe *probe.Probe } // NewReporter makes a new Reporter -func NewReporter(registry Registry, hostID string) *Reporter { +func NewReporter(registry Registry, hostID string, probe *probe.Probe) *Reporter { return &Reporter{ registry: registry, hostID: hostID, + probe: probe, } } +// ContainerUpdated should be called whenever a container is updated. +func (r *Reporter) ContainerUpdated(c Container) { + localAddrs, err := report.LocalAddresses() + if err != nil { + log.Printf("Error getting local address: %v", err) + return + } + + // Publish a 'short cut' report container just this container + rpt := report.MakeReport() + rpt.Container.AddNode(report.MakeContainerNodeID(r.hostID, c.ID()), c.GetNode(r.hostID, localAddrs)) + r.probe.Publish(rpt) +} + // Report generates a Report containing Container and ContainerImage topologies func (r *Reporter) Report() (report.Report, error) { localAddrs, err := report.LocalAddresses() diff --git a/probe/docker/reporter_test.go b/probe/docker/reporter_test.go index cbabba81b..598773b69 100644 --- a/probe/docker/reporter_test.go +++ b/probe/docker/reporter_test.go @@ -36,6 +36,8 @@ func (r *mockRegistry) WalkImages(f func(*client.APIImages)) { } } +func (r *mockRegistry) WatchContainerUpdates(_ docker.ContainerUpdateWatcher) {} + var ( mockRegistryInstance = &mockRegistry{ containersByPID: map[int]docker.Container{ @@ -95,7 +97,7 @@ func TestReporter(t *testing.T) { Controls: report.Controls{}, } - reporter := docker.NewReporter(mockRegistryInstance, "") + reporter := docker.NewReporter(mockRegistryInstance, "", nil) have, _ := reporter.Report() if !reflect.DeepEqual(want, have) { t.Errorf("%s", test.Diff(want, have)) diff --git a/probe/probe.go b/probe/probe.go index d429e7cdf..57edba645 100644 --- a/probe/probe.go +++ b/probe/probe.go @@ -9,6 +9,10 @@ import ( "github.com/weaveworks/scope/xfer" ) +const ( + reportBufferSize = 16 +) + // Probe sits there, generating and publishing reports. type Probe struct { spyInterval, publishInterval time.Duration @@ -20,7 +24,9 @@ type Probe struct { quit chan struct{} done sync.WaitGroup - rpt syncReport + + spiedReports chan report.Report + shortcutReports chan report.Report } // Tagger tags nodes with value-add node metadata. @@ -47,8 +53,9 @@ func New(spyInterval, publishInterval time.Duration, publisher xfer.Publisher) * publishInterval: publishInterval, publisher: publisher, quit: make(chan struct{}), + spiedReports: make(chan report.Report, reportBufferSize), + shortcutReports: make(chan report.Report, reportBufferSize), } - result.rpt.swap(report.MakeReport()) return result } @@ -80,6 +87,12 @@ func (p *Probe) Stop() { p.done.Wait() } +// Publish will queue a report for immediate publication, +// bypassing the spy tick +func (p *Probe) Publish(rpt report.Report) { + p.shortcutReports <- rpt +} + func (p *Probe) spyLoop() { defer p.done.Done() spyTick := time.Tick(p.spyInterval) @@ -94,10 +107,9 @@ func (p *Probe) spyLoop() { } } - localReport := p.rpt.copy() - localReport = localReport.Merge(p.report()) - localReport = p.tag(localReport) - p.rpt.swap(localReport) + rpt := p.report() + rpt = p.tag(rpt) + p.spiedReports <- rpt if took := time.Since(start); took > p.spyInterval { log.Printf("report generation took too long (%s)", took) @@ -140,6 +152,17 @@ func (p *Probe) tag(r report.Report) report.Report { return r } +func condense(rpt report.Report, rs chan report.Report) report.Report { + for { + select { + case r := <-rs: + rpt = rpt.Merge(r) + default: + return rpt + } + } +} + func (p *Probe) publishLoop() { defer p.done.Done() var ( @@ -150,8 +173,14 @@ func (p *Probe) publishLoop() { for { select { case <-pubTick: - localReport := p.rpt.swap(report.MakeReport()) - if err := rptPub.Publish(localReport); err != nil { + rpt := condense(report.MakeReport(), p.spiedReports) + if err := rptPub.Publish(rpt); err != nil { + log.Printf("publish: %v", err) + } + + case rpt := <-p.shortcutReports: + rpt = condense(rpt, p.shortcutReports) + if err := rptPub.Publish(rpt); err != nil { log.Printf("publish: %v", err) } diff --git a/probe/sync_report.go b/probe/sync_report.go deleted file mode 100644 index 830214920..000000000 --- a/probe/sync_report.go +++ /dev/null @@ -1,26 +0,0 @@ -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() -} diff --git a/prog/probe/main.go b/prog/probe/main.go index 4b4e19400..536e0e572 100644 --- a/prog/probe/main.go +++ b/prog/probe/main.go @@ -134,7 +134,10 @@ func main() { if registry, err := docker.NewRegistry(*dockerInterval); err == nil { defer registry.Stop() p.AddTagger(docker.NewTagger(registry, processCache)) - p.AddReporter(docker.NewReporter(registry, hostID)) + + reporter := docker.NewReporter(registry, hostID, p) + registry.WatchContainerUpdates(reporter.ContainerUpdated) + p.AddReporter(reporter) } else { log.Printf("Docker: failed to start registry: %v", err) } From 47ef1e02b3fd4c6cdb7ca383c1424a0c051f3573 Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Tue, 10 Nov 2015 11:16:06 +0000 Subject: [PATCH 2/3] Shortcut app -> UI ws push for certain reports. --- app/api_topology.go | 7 ++++++- app/collector.go | 38 ++++++++++++++++++++++++++++++++++++++ app/mock_reporter_test.go | 5 +++-- probe/docker/reporter.go | 1 + report/report.go | 4 ++++ 5 files changed, 52 insertions(+), 3 deletions(-) diff --git a/app/api_topology.go b/app/api_topology.go index 315f7ce8e..00d86c85a 100644 --- a/app/api_topology.go +++ b/app/api_topology.go @@ -114,7 +114,11 @@ func handleWebsocket( var ( previousTopo render.RenderableNodes tick = time.Tick(loop) + wait = make(chan struct{}, 1) ) + rep.WaitOn(wait) + defer rep.UnWait(wait) + for { newTopo := renderer.Render(rep.Report()).Prune() diff := render.TopoDiff(previousTopo, newTopo) @@ -128,9 +132,10 @@ func handleWebsocket( } select { + case <-wait: + case <-tick: case <-quit: return - case <-tick: } } } diff --git a/app/collector.go b/app/collector.go index d97846ab2..394043c52 100644 --- a/app/collector.go +++ b/app/collector.go @@ -11,6 +11,8 @@ import ( // interface for parts of the app, and several experimental components. type Reporter interface { Report() report.Report + WaitOn(chan struct{}) + UnWait(chan struct{}) } // Adder is something that can accept reports. It's a convenient interface for @@ -25,12 +27,45 @@ type Collector struct { mtx sync.Mutex reports []timestampReport window time.Duration + waitableCondition +} + +type waitableCondition struct { + sync.Mutex + waiters map[chan struct{}]struct{} +} + +func (wc *waitableCondition) WaitOn(waiter chan struct{}) { + wc.Lock() + wc.waiters[waiter] = struct{}{} + wc.Unlock() +} + +func (wc *waitableCondition) UnWait(waiter chan struct{}) { + wc.Lock() + delete(wc.waiters, waiter) + wc.Unlock() +} + +func (wc *waitableCondition) Broadcast() { + wc.Lock() + for waiter := range wc.waiters { + // Non-block write to channel + select { + case waiter <- struct{}{}: + default: + } + } + wc.Unlock() } // NewCollector returns a collector ready for use. func NewCollector(window time.Duration) *Collector { return &Collector{ window: window, + waitableCondition: waitableCondition{ + waiters: map[chan struct{}]struct{}{}, + }, } } @@ -42,6 +77,9 @@ func (c *Collector) Add(rpt report.Report) { defer c.mtx.Unlock() c.reports = append(c.reports, timestampReport{now(), rpt}) c.reports = clean(c.reports, c.window) + if rpt.Shortcut { + c.Broadcast() + } } // Report returns a merged report over all added reports. It implements diff --git a/app/mock_reporter_test.go b/app/mock_reporter_test.go index 99439e9a0..9e27941c1 100644 --- a/app/mock_reporter_test.go +++ b/app/mock_reporter_test.go @@ -9,5 +9,6 @@ import ( type StaticReport struct{} func (s StaticReport) Report() report.Report { return fixture.Report } - -func (s StaticReport) Add(report.Report) {} +func (s StaticReport) Add(report.Report) {} +func (s StaticReport) WaitOn(chan struct{}) {} +func (s StaticReport) UnWait(chan struct{}) {} diff --git a/probe/docker/reporter.go b/probe/docker/reporter.go index a2c10c19e..afd8e565a 100644 --- a/probe/docker/reporter.go +++ b/probe/docker/reporter.go @@ -42,6 +42,7 @@ func (r *Reporter) ContainerUpdated(c Container) { // Publish a 'short cut' report container just this container rpt := report.MakeReport() + rpt.Shortcut = true rpt.Container.AddNode(report.MakeContainerNodeID(r.hostID, c.ID()), c.GetNode(r.hostID, localAddrs)) r.probe.Publish(rpt) } diff --git a/report/report.go b/report/report.go index aca0ad2fd..88876ea6d 100644 --- a/report/report.go +++ b/report/report.go @@ -63,6 +63,10 @@ type Report struct { // such as in the app, we expect the component to overwrite the window // before serving it to consumers. Window time.Duration + + // Shortcut reports should be propogated to the UI as quickly as possible, + // bypassing the usual spy interval, publish interval and app ws interval. + Shortcut bool } // MakeReport makes a clean report, ready to Merge() other reports into. From c610bf0ea1b3859a2b45780a32c805206fa8dc64 Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Wed, 11 Nov 2015 11:42:28 +0000 Subject: [PATCH 3/3] Review feedback. --- probe/docker/reporter.go | 4 +++- probe/probe.go | 28 ++++++++++++---------------- prog/probe/main.go | 5 +---- 3 files changed, 16 insertions(+), 21 deletions(-) diff --git a/probe/docker/reporter.go b/probe/docker/reporter.go index afd8e565a..6d1e088da 100644 --- a/probe/docker/reporter.go +++ b/probe/docker/reporter.go @@ -25,11 +25,13 @@ type Reporter struct { // NewReporter makes a new Reporter func NewReporter(registry Registry, hostID string, probe *probe.Probe) *Reporter { - return &Reporter{ + reporter := &Reporter{ registry: registry, hostID: hostID, probe: probe, } + registry.WatchContainerUpdates(reporter.ContainerUpdated) + return reporter } // ContainerUpdated should be called whenever a container is updated. diff --git a/probe/probe.go b/probe/probe.go index 57edba645..a0700ba2f 100644 --- a/probe/probe.go +++ b/probe/probe.go @@ -16,7 +16,7 @@ const ( // Probe sits there, generating and publishing reports. type Probe struct { spyInterval, publishInterval time.Duration - publisher xfer.Publisher + publisher *xfer.ReportPublisher tickers []Ticker reporters []Reporter @@ -51,7 +51,7 @@ func New(spyInterval, publishInterval time.Duration, publisher xfer.Publisher) * result := &Probe{ spyInterval: spyInterval, publishInterval: publishInterval, - publisher: publisher, + publisher: xfer.NewReportPublisher(publisher), quit: make(chan struct{}), spiedReports: make(chan report.Report, reportBufferSize), shortcutReports: make(chan report.Report, reportBufferSize), @@ -152,37 +152,33 @@ func (p *Probe) tag(r report.Report) report.Report { return r } -func condense(rpt report.Report, rs chan report.Report) report.Report { +func (p *Probe) drainAndPublish(rpt report.Report, rs chan report.Report) { +ForLoop: for { select { case r := <-rs: rpt = rpt.Merge(r) default: - return rpt + break ForLoop } } + + if err := p.publisher.Publish(rpt); err != nil { + log.Printf("publish: %v", err) + } } func (p *Probe) publishLoop() { defer p.done.Done() - var ( - pubTick = time.Tick(p.publishInterval) - rptPub = xfer.NewReportPublisher(p.publisher) - ) + pubTick := time.Tick(p.publishInterval) for { select { case <-pubTick: - rpt := condense(report.MakeReport(), p.spiedReports) - if err := rptPub.Publish(rpt); err != nil { - log.Printf("publish: %v", err) - } + p.drainAndPublish(report.MakeReport(), p.spiedReports) case rpt := <-p.shortcutReports: - rpt = condense(rpt, p.shortcutReports) - if err := rptPub.Publish(rpt); err != nil { - log.Printf("publish: %v", err) - } + p.drainAndPublish(rpt, p.shortcutReports) case <-p.quit: return diff --git a/prog/probe/main.go b/prog/probe/main.go index 536e0e572..6e2d24222 100644 --- a/prog/probe/main.go +++ b/prog/probe/main.go @@ -134,10 +134,7 @@ func main() { if registry, err := docker.NewRegistry(*dockerInterval); err == nil { defer registry.Stop() p.AddTagger(docker.NewTagger(registry, processCache)) - - reporter := docker.NewReporter(registry, hostID, p) - registry.WatchContainerUpdates(reporter.ContainerUpdated) - p.AddReporter(reporter) + p.AddReporter(docker.NewReporter(registry, hostID, p)) } else { log.Printf("Docker: failed to start registry: %v", err) }