From 092167f5de739c7581a5f8b969c10f8aa02b0a92 Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Thu, 1 Oct 2015 10:44:44 +0000 Subject: [PATCH] Don't allow lagging report generation to prevent report publication. --- probe/main.go | 75 +++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 55 insertions(+), 20 deletions(-) diff --git a/probe/main.go b/probe/main.go index e14beddd1..eacbc2cf9 100644 --- a/probe/main.go +++ b/probe/main.go @@ -10,6 +10,7 @@ import ( "os" "os/signal" "strings" + "sync" "syscall" "time" @@ -153,27 +154,22 @@ func main() { }() } - quit, done := make(chan struct{}), make(chan struct{}) - defer func() { <-done }() // second, wait for the main loop to be killed - defer close(quit) // first, kill the main loop + quit, done := make(chan struct{}), sync.WaitGroup{} + done.Add(2) + defer func() { done.Wait() }() // second, wait for the main loops to be killed + defer close(quit) // first, kill the main loops + + var ( + rpt = report.MakeReport() + rptLock = sync.Mutex{} + ) + go func() { - defer close(done) - var ( - pubTick = time.Tick(*publishInterval) - spyTick = time.Tick(*spyInterval) - r = report.MakeReport() - p = xfer.NewReportPublisher(publishers) - ) + defer done.Done() + spyTick := time.Tick(*spyInterval) + for { select { - case <-pubTick: - publishTicks.WithLabelValues().Add(1) - r.Window = *publishInterval - if err := p.Publish(r); err != nil { - log.Printf("publish: %v", err) - } - r = report.MakeReport() - case <-spyTick: start := time.Now() for _, ticker := range tickers { @@ -181,8 +177,18 @@ func main() { log.Printf("error doing ticker: %v", err) } } - r = r.Merge(doReport(reporters)) - r = Apply(r, taggers) + + rptLock.Lock() + localReport := rpt.Copy() + rptLock.Unlock() + + localReport = localReport.Merge(doReport(reporters)) + localReport = Apply(localReport, taggers) + + rptLock.Lock() + rpt = localReport + rptLock.Unlock() + if took := time.Since(start); took > *spyInterval { log.Printf("report generation took too long (%s)", took) } @@ -193,6 +199,35 @@ func main() { } }() + go func() { + defer done.Done() + var ( + pubTick = time.Tick(*publishInterval) + p = xfer.NewReportPublisher(publishers) + localReport = report.MakeReport() + ) + + for { + select { + case <-pubTick: + publishTicks.WithLabelValues().Add(1) + + rptLock.Lock() + localReport = rpt + rpt = report.MakeReport() + rptLock.Unlock() + + localReport.Window = *publishInterval + if err := p.Publish(localReport); err != nil { + log.Printf("publish: %v", err) + } + + case <-quit: + return + } + } + }() + log.Printf("%s", <-interrupt()) }