From 1eb57c2e40c0daa8c8b841ea8820bc8fca41f5c7 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sat, 27 Mar 2021 16:11:03 +0000 Subject: [PATCH 01/11] Multitenant collector now always saves async Removed support for saving all reports immediately --- app/multitenant/aws_collector.go | 29 +++++++++++++---------------- 1 file changed, 13 insertions(+), 16 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 03c7203d7..5f141e938 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -205,6 +205,7 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { waiters: map[watchKey]*nats.Subscription{}, } + // If given a StoreInterval we will be storing periodically; if not we only answer queries if config.StoreInterval != 0 { c.ticker = time.NewTicker(config.StoreInterval) go c.flushLoop() @@ -681,6 +682,9 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e if err != nil { return err } + if c.cfg.StoreInterval == 0 { + return fmt.Errorf("--app.collector.store-interval must be non-zero") + } // Shortcut reports are published to nats but not persisted - // we'll get a full report from the same probe in a few seconds @@ -702,23 +706,16 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e return nil } - if c.cfg.StoreInterval == 0 { - rowKey, colKey, reportKey := calculateReportKeys(userid, time.Now()) - err = c.persistReport(ctx, userid, rowKey, colKey, reportKey, buf) - if err != nil { - return err - } - } else { - rep = c.massageReport(userid, rep) - entry := &pendingEntry{report: report.MakeReport()} - if e, found := c.pending.LoadOrStore(userid, entry); found { - entry = e.(*pendingEntry) - } - entry.Lock() - entry.report.UnsafeMerge(rep) - entry.count++ - entry.Unlock() + // We are building up a report in memory; merge into that and it will be saved shortly + rep = c.massageReport(userid, rep) + entry := &pendingEntry{report: report.MakeReport()} + if e, found := c.pending.LoadOrStore(userid, entry); found { + entry = e.(*pendingEntry) } + entry.Lock() + entry.report.UnsafeMerge(rep) + entry.count++ + entry.Unlock() return nil } From b9c8cf6998302020d59a139de5c1ed8c2cc09baa Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 28 Mar 2021 13:59:25 +0100 Subject: [PATCH 02/11] Add flag for querier to talk to collectors --- app/multitenant/aws_collector.go | 1 + prog/app.go | 5 +++-- prog/main.go | 4 +++- 3 files changed, 7 insertions(+), 3 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 5f141e938..868c0a91a 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -142,6 +142,7 @@ type AWSCollectorConfig struct { MemcacheClient *MemcacheClient Window time.Duration MaxTopNodes int + CollectorAddr string } // if StoreInterval is set, reports are merged into here and held until flushed to store diff --git a/prog/app.go b/prog/app.go index 5648b57d9..dff8a88da 100644 --- a/prog/app.go +++ b/prog/app.go @@ -90,7 +90,7 @@ func router(collector app.Collector, controlRouter app.ControlRouter, pipeRouter } func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL string, storeInterval time.Duration, natsHostname string, - memcacheConfig multitenant.MemcacheConfig, window time.Duration, maxTopNodes int, createTables bool) (app.Collector, error) { + memcacheConfig multitenant.MemcacheConfig, window time.Duration, maxTopNodes int, createTables bool, collectorAddr string) (app.Collector, error) { if collectorURL == "local" { return app.NewCollector(window), nil } @@ -134,6 +134,7 @@ func collectorFactory(userIDer multitenant.UserIDer, collectorURL, s3URL string, MemcacheClient: memcacheClient, Window: window, MaxTopNodes: maxTopNodes, + CollectorAddr: collectorAddr, }, ) if err != nil { @@ -248,7 +249,7 @@ func appMain(flags appFlags) { Service: flags.memcachedService, CompressionLevel: flags.memcachedCompressionLevel, }, - flags.window, flags.maxTopNodes, flags.awsCreateTables) + flags.window, flags.maxTopNodes, flags.awsCreateTables, flags.collectorAddr) if err != nil { log.Fatalf("Error creating collector: %v", err) return diff --git a/prog/main.go b/prog/main.go index 30e9e1f86..66eb5c4c7 100644 --- a/prog/main.go +++ b/prog/main.go @@ -162,7 +162,8 @@ type appFlags struct { containerName string dockerEndpoint string - collectorURL string + collectorURL string // how collector talks to backing store (or "local" if none) + collectorAddr string // how to find collectors if deployed as microservices s3URL string storeInterval time.Duration controlRouterURL string @@ -376,6 +377,7 @@ func setupFlags(flags *flags) { flag.Var(&flags.containerLabelFilterFlagsExclude, "app.container-label-filter-exclude", "Add container label-based view filter that excludes containers with the given label, specified as title:label. Multiple flags are accepted. Example: --app.container-label-filter-exclude='Database Containers:role=db'") flag.StringVar(&flags.app.collectorURL, "app.collector", "local", "Collector to use (local, dynamodb, or file/directory)") + flag.StringVar(&flags.app.collectorAddr, "app.collector-addr", "", "Address to look up collectors when deployed as microservices") flag.StringVar(&flags.app.s3URL, "app.collector.s3", "local", "S3 URL to use (when collector is dynamodb)") flag.DurationVar(&flags.app.storeInterval, "app.collector.store-interval", 0, "How often to store merged incoming reports. If 0, reports are stored unmerged as they arrive.") flag.StringVar(&flags.app.controlRouterURL, "app.control.router", "local", "Control router to use (local or sqs)") From 5d12b7ff6534a83822fae4821dca0afb34592cbf Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 28 Mar 2021 14:06:21 +0100 Subject: [PATCH 03/11] Refactor: extract multitenant collection of 'live' reports To help clarify subsequent changes --- app/multitenant/aws_collector.go | 17 +--------------- app/multitenant/collector.go | 34 ++++++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 16 deletions(-) create mode 100644 app/multitenant/collector.go diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 868c0a91a..28273ebf5 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -145,13 +145,6 @@ type AWSCollectorConfig struct { CollectorAddr string } -// if StoreInterval is set, reports are merged into here and held until flushed to store -type pendingEntry struct { - sync.Mutex - report report.Report - count int -} - type awsCollector struct { cfg AWSCollectorConfig db *dynamodb.DynamoDB @@ -707,16 +700,8 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) e return nil } - // We are building up a report in memory; merge into that and it will be saved shortly rep = c.massageReport(userid, rep) - entry := &pendingEntry{report: report.MakeReport()} - if e, found := c.pending.LoadOrStore(userid, entry); found { - entry = e.(*pendingEntry) - } - entry.Lock() - entry.report.UnsafeMerge(rep) - entry.count++ - entry.Unlock() + c.addToLive(ctx, userid, rep) return nil } diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go new file mode 100644 index 000000000..ccbf14ed2 --- /dev/null +++ b/app/multitenant/collector.go @@ -0,0 +1,34 @@ +package multitenant + +// Collect reports from probes per-tenant, and supply them to queriers on demand + +import ( + "sync" + + "context" + + "github.com/weaveworks/scope/report" +) + +// if StoreInterval is set, reports are merged into here and held until flushed to store +type pendingEntry struct { + sync.Mutex + report *report.Report +} + +// We are building up a report in memory; merge into that and it will be saved shortly +// NOTE: may retain a reference to rep; must not be used by caller after this. +func (c *awsCollector) addToLive(ctx context.Context, userid string, rep report.Report) { + entry := &pendingEntry{} + if e, found := c.pending.LoadOrStore(userid, entry); found { + entry = e.(*pendingEntry) + } + entry.Lock() + if entry.report == nil { + entry.report = &rep + } else { + entry.report.UnsafeMerge(rep) + } + entry.Unlock() +} + From 667daef81b555a2c77acbe669779a44bea587723 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 28 Mar 2021 14:09:00 +0100 Subject: [PATCH 04/11] Refactor: extract function reportsFromStore() To help clarify subsequent changes --- app/multitenant/aws_collector.go | 39 ++++++++++++++++++++------------ 1 file changed, 25 insertions(+), 14 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 28273ebf5..ba9853627 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -459,15 +459,6 @@ func (c *awsCollector) massageReport(userid string, report report.Report) report return report } -/* -S3 stores original reports from one probe at the timestamp they arrived at collector. -Collector also sends every report to memcached. -The in-memory cache stores: - - individual reports deserialised, under S3 key for report - - sets of reports in interval [t,t+3) merged, under key "instance:t" - - so to check the cache for reports from 14:31:00 to 14:31:15 you would request 5 keys 3 seconds apart -*/ - func (c *awsCollector) Report(ctx context.Context, timestamp time.Time) (report.Report, error) { span, ctx := opentracing.StartSpanFromContext(ctx, "awsCollector.Report") defer span.Finish() @@ -476,11 +467,32 @@ func (c *awsCollector) Report(ctx context.Context, timestamp time.Time) (report. return report.MakeReport(), err } span.SetTag("userid", userid) + var reports []report.Report + reports, err = c.reportsFromStore(ctx, userid, timestamp) + if err != nil { + return report.MakeReport(), err + } + span.LogFields(otlog.Int("merging", len(reports))) + return c.merger.Merge(reports), nil +} + +/* +Given a timestamp in the past, fetch reports within the window from store or cache + +S3 stores original reports from one probe at the timestamp they arrived at collector. +Collector also sends every report to memcached. +The in-memory cache stores: + - individual reports deserialised, under S3 key for report + - sets of reports in interval [t,t+3) merged, under key "instance:t" + - so to check the cache for reports from 14:31:00 to 14:31:15 you would request 5 keys 3 seconds apart +*/ +func (c *awsCollector) reportsFromStore(ctx context.Context, userid string, timestamp time.Time) ([]report.Report, error) { + span := opentracing.SpanFromContext(ctx) end := timestamp start := end.Add(-c.cfg.Window) reportKeys, err := c.getReportKeys(ctx, userid, start, end) if err != nil { - return report.MakeReport(), err + return nil, err } span.LogFields(otlog.Int("keys", len(reportKeys)), otlog.String("timestamp", timestamp.String())) @@ -491,18 +503,17 @@ func (c *awsCollector) Report(ctx context.Context, timestamp time.Time) (report. for ; ts+(reportQuantisationInterval+gracePeriod).Nanoseconds() < endTS; ts += reportQuantisationInterval.Nanoseconds() { quantumReport, err := c.reportForQuantum(ctx, userid, reportKeys, ts) if err != nil { - return report.MakeReport(), err + return nil, err } reports = append(reports, quantumReport) } // Fetch individual reports for the period after the last quantum last, err := c.reportsForKeysInRange(ctx, userid, reportKeys, ts, endTS) if err != nil { - return report.MakeReport(), err + return nil, err } reports = append(reports, last...) - span.LogFields(otlog.Int("merging", len(reports))) - return c.merger.Merge(reports), nil + return reports, nil } // Fetch a merged report either from cache or from store which we put in cache From 5032cca5c0c6c702570895e735fe6a386b3723bb Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 28 Mar 2021 14:11:26 +0100 Subject: [PATCH 05/11] Multitenant mode: fetch live data from collectors Collectors hold recent reports in memory. When querier needs 'live' data, fetch it from collectors instead of from the long-term store. Send reports from collector to querier in msgpack; disable compression on REST call, otherwise Go silently decompresses, which takes longer. --- app/api_report.go | 2 +- app/multitenant/aws_collector.go | 20 +++++-- app/multitenant/collector.go | 100 +++++++++++++++++++++++++++++++ app/server_helpers.go | 24 ++++++++ go.mod | 1 + vendor/modules.txt | 1 + 6 files changed, 143 insertions(+), 5 deletions(-) diff --git a/app/api_report.go b/app/api_report.go index 7ab162528..079625b54 100644 --- a/app/api_report.go +++ b/app/api_report.go @@ -19,7 +19,7 @@ func makeRawReportHandler(rep Reporter) CtxHandlerFunc { return } censorCfg := report.GetCensorConfigFromRequest(r) - respondWith(ctx, w, http.StatusOK, report.CensorRawReport(rawReport, censorCfg)) + respondWithReport(ctx, w, r, report.CensorRawReport(rawReport, censorCfg)) } } diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index ba9853627..551db2fcc 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -242,11 +242,17 @@ func (c *awsCollector) flushPending(ctx context.Context) { entry := value.(*pendingEntry) entry.Lock() - rpt, count := entry.report, entry.count - entry.report, entry.count = report.MakeReport(), 0 + rpt := entry.report + entry.report = nil + if entry.older == nil { + entry.older = make([]*report.Report, c.cfg.Window/c.cfg.StoreInterval) + } else { + copy(entry.older[1:], entry.older) // move everything down one + } + entry.older[0] = rpt entry.Unlock() - if count > 0 { + if rpt != nil { // serialise reports on one goroutine to limit CPU usage buf, err := rpt.WriteBinary() if err != nil { @@ -459,6 +465,8 @@ func (c *awsCollector) massageReport(userid string, report report.Report) report return report } +// If we are running as a Query service, fetch data and merge into a report +// If we are running as a Collector and the request is for live data, merge in-memory data and return func (c *awsCollector) Report(ctx context.Context, timestamp time.Time) (report.Report, error) { span, ctx := opentracing.StartSpanFromContext(ctx, "awsCollector.Report") defer span.Finish() @@ -468,7 +476,11 @@ func (c *awsCollector) Report(ctx context.Context, timestamp time.Time) (report. } span.SetTag("userid", userid) var reports []report.Report - reports, err = c.reportsFromStore(ctx, userid, timestamp) + if time.Since(timestamp) < c.cfg.Window { + reports, err = c.reportsFromLive(ctx, userid) + } else { + reports, err = c.reportsFromStore(ctx, userid, timestamp) + } if err != nil { return report.MakeReport(), err } diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go index ccbf14ed2..9be4668e6 100644 --- a/app/multitenant/collector.go +++ b/app/multitenant/collector.go @@ -3,10 +3,19 @@ package multitenant // Collect reports from probes per-tenant, and supply them to queriers on demand import ( + "fmt" + "io" + "net" + "net/http" + "strconv" "sync" "context" + "github.com/opentracing-contrib/go-stdlib/nethttp" + opentracing "github.com/opentracing/opentracing-go" + log "github.com/sirupsen/logrus" + "github.com/weaveworks/common/user" "github.com/weaveworks/scope/report" ) @@ -14,6 +23,7 @@ import ( type pendingEntry struct { sync.Mutex report *report.Report + older []*report.Report } // We are building up a report in memory; merge into that and it will be saved shortly @@ -32,3 +42,93 @@ func (c *awsCollector) addToLive(ctx context.Context, userid string, rep report. entry.Unlock() } +func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]report.Report, error) { + span, ctx := opentracing.StartSpanFromContext(ctx, "reportsFromLive") + defer span.Finish() + if c.cfg.StoreInterval != 0 { + // We are a collector + e, found := c.pending.Load(userid) + if !found { + return nil, nil + } + entry := e.(*pendingEntry) + entry.Lock() + ret := make([]report.Report, 0, len(entry.older)+1) + if entry.report != nil { + ret = append(ret, entry.report.Copy()) // Copy contents because this report is being unsafe-merged to + } + for _, v := range entry.older { + if v != nil { + ret = append(ret, *v) // no copy because older reports are immutable + } + } + entry.Unlock() + return ret, nil + } + + // We are a querier: fetch the most up-to-date reports from collectors + // TODO: resolve c.collectorAddress periodically instead of every time we make a call + addrs := resolve(c.cfg.CollectorAddr) + ret := make([]report.Report, 0, len(addrs)) + // make a call to each collector and fetch its data for this userid + // TODO: do them in parallel + for _, addr := range addrs { + body, err := oneCall(ctx, addr, "/api/report", userid) + if err != nil { + log.Warnf("error calling '%s': %v", addr, err) + continue + } + rpt, err := report.MakeFromBinary(ctx, body, false, true) + body.Close() + if err != nil { + log.Warnf("error decoding: %v", err) + continue + } + ret = append(ret, *rpt) + } + + return ret, nil +} + +func resolve(name string) []string { + _, addrs, err := net.LookupSRV("", "", name) + if err != nil { + log.Warnf("Cannot resolve '%s': %v", name, err) + return []string{} + } + endpoints := make([]string, 0, len(addrs)) + for _, addr := range addrs { + port := strconv.Itoa(int(addr.Port)) + endpoints = append(endpoints, net.JoinHostPort(addr.Target, port)) + } + return endpoints +} + +func oneCall(ctx context.Context, endpoint, path, userid string) (io.ReadCloser, error) { + fullPath := "http://" + endpoint + path + req, err := http.NewRequest("GET", fullPath, nil) + if err != nil { + return nil, fmt.Errorf("error making request %s: %w", fullPath, err) + } + req = req.WithContext(ctx) + req.Header.Set(user.OrgIDHeaderName, userid) + req.Header.Set("Accept", "application/msgpack") + req.Header.Set("Accept-Encoding", "identity") // disable compression + if parentSpan := opentracing.SpanFromContext(ctx); parentSpan != nil { + var ht *nethttp.Tracer + req, ht = nethttp.TraceRequest(parentSpan.Tracer(), req, nethttp.OperationName("Collector Fetch")) + defer ht.Finish() + } + client := &http.Client{Transport: &nethttp.Transport{}} + res, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("error getting %s: %w", fullPath, err) + } + if res.StatusCode != http.StatusOK { + content, _ := io.ReadAll(res.Body) + res.Body.Close() + return nil, fmt.Errorf("error from collector: %s (%s)", res.Status, string(content)) + } + + return res.Body, nil +} diff --git a/app/server_helpers.go b/app/server_helpers.go index a9970c704..42d19fdfb 100644 --- a/app/server_helpers.go +++ b/app/server_helpers.go @@ -3,6 +3,7 @@ package app import ( "context" "net/http" + "strings" opentracing "github.com/opentracing/opentracing-go" "github.com/ugorji/go/codec" @@ -33,3 +34,26 @@ func respondWith(ctx context.Context, w http.ResponseWriter, code int, response log.Errorf("Error encoding response: %v", err) } } + +// Similar to the above function, but respect the request's Accept header. +// Possibly we should do a complete parse of Accept, but for now just rudimentary check +func respondWithReport(ctx context.Context, w http.ResponseWriter, req *http.Request, response interface{}) { + accept := req.Header.Get("Accept") + if strings.HasPrefix(accept, "application/msgpack") { + w.Header().Set("Content-Type", "application/msgpack") + w.WriteHeader(http.StatusOK) + encoder := codec.NewEncoder(w, &codec.MsgpackHandle{}) + if err := encoder.Encode(response); err != nil { + log.Errorf("Error encoding response: %v", err) + } + return + } + + w.Header().Set("Content-Type", "application/json") + w.Header().Add("Cache-Control", "no-cache") + w.WriteHeader(http.StatusOK) + encoder := codec.NewEncoder(w, &codec.JsonHandle{}) + if err := encoder.Encode(response); err != nil { + log.Errorf("Error encoding response: %v", err) + } +} diff --git a/go.mod b/go.mod index ae4052136..73322e315 100644 --- a/go.mod +++ b/go.mod @@ -52,6 +52,7 @@ require ( github.com/nats-io/nuid v0.0.0-20160402145409-a5152d67cf63 // indirect github.com/opencontainers/runc v1.0.0-rc5 // indirect github.com/openebs/k8s-snapshot-client v0.0.0-20180831100134-a6506305fb16 + github.com/opentracing-contrib/go-stdlib v0.0.0-20190519235532-cf7a6c988dc9 github.com/opentracing/opentracing-go v1.1.0 github.com/paypal/ionet v0.0.0-20130919195445-ed0aaebc5417 github.com/pborman/uuid v0.0.0-20150824212802-cccd189d45f7 diff --git a/vendor/modules.txt b/vendor/modules.txt index 1c1a3c808..2410cb25b 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -293,6 +293,7 @@ github.com/openebs/k8s-snapshot-client/snapshot/pkg/client/clientset/versioned github.com/openebs/k8s-snapshot-client/snapshot/pkg/client/clientset/versioned/scheme github.com/openebs/k8s-snapshot-client/snapshot/pkg/client/clientset/versioned/typed/volumesnapshot/v1 # github.com/opentracing-contrib/go-stdlib v0.0.0-20190519235532-cf7a6c988dc9 +## explicit github.com/opentracing-contrib/go-stdlib/nethttp # github.com/opentracing/opentracing-go v1.1.0 ## explicit From bea8db3683fd5dcfffea7325d85d6e85ac530307 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 11 Apr 2021 12:23:29 +0000 Subject: [PATCH 06/11] Log/trace data size before decoding report This lets us see when the reading finished and decoding started --- report/marshal.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/report/marshal.go b/report/marshal.go index 38d98253d..c93b71071 100644 --- a/report/marshal.go +++ b/report/marshal.go @@ -102,10 +102,6 @@ func MakeFromBinary(ctx context.Context, r io.Reader, gzipped bool, msgpack bool if err != nil { return nil, err } - rep := MakeReport() - if err := codec.NewDecoderBytes(buf.Bytes(), codecHandle(msgpack)).Decode(&rep); err != nil { - return nil, err - } log.Debugf( "Received report sizes: compressed %d bytes, uncompressed %d bytes (%.2f%%)", compressedSize, @@ -113,6 +109,10 @@ func MakeFromBinary(ctx context.Context, r io.Reader, gzipped bool, msgpack bool float32(compressedSize)/float32(uncompressedSize)*100, ) span.LogFields(otlog.Uint64("compressedSize", compressedSize), otlog.Int64("uncompressedSize", uncompressedSize)) + rep := MakeReport() + if err := codec.NewDecoderBytes(buf.Bytes(), codecHandle(msgpack)).Decode(&rep); err != nil { + return nil, err + } return &rep, nil } From 055ca53241594dc8fcbbf916de426622588618d7 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 11 Apr 2021 13:03:42 +0000 Subject: [PATCH 07/11] refactor: extract fn to check whether collector or querier --- app/multitenant/aws_collector.go | 2 +- app/multitenant/collector.go | 7 +++++-- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 551db2fcc..0502e0321 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -200,7 +200,7 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { } // If given a StoreInterval we will be storing periodically; if not we only answer queries - if config.StoreInterval != 0 { + if c.isCollector() { c.ticker = time.NewTicker(config.StoreInterval) go c.flushLoop() } diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go index 9be4668e6..b6bcc7107 100644 --- a/app/multitenant/collector.go +++ b/app/multitenant/collector.go @@ -42,11 +42,14 @@ func (c *awsCollector) addToLive(ctx context.Context, userid string, rep report. entry.Unlock() } +func (c *awsCollector) isCollector() bool { + return c.cfg.StoreInterval != 0 +} + func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]report.Report, error) { span, ctx := opentracing.StartSpanFromContext(ctx, "reportsFromLive") defer span.Finish() - if c.cfg.StoreInterval != 0 { - // We are a collector + if c.isCollector() { e, found := c.pending.Load(userid) if !found { return nil, nil From 99582ba8357a71870ea481a8ef99bde43c375f1c Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 12 Apr 2021 19:44:03 +0000 Subject: [PATCH 08/11] Implement HasReports for live data from collectors --- app/multitenant/aws_collector.go | 4 +++ app/multitenant/collector.go | 43 ++++++++++++++++++++++++++++++++ 2 files changed, 47 insertions(+) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 0502e0321..0a387e5c1 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -564,6 +564,10 @@ func (c *awsCollector) HasReports(ctx context.Context, timestamp time.Time) (boo if err != nil { return false, err } + if time.Since(timestamp) < c.cfg.Window { + has, err := c.hasReportsFromLive(ctx, userid) + return has, err + } start := timestamp.Add(-c.cfg.Window) reportKeys, err := c.getReportKeys(ctx, userid, start, timestamp) return len(reportKeys) > 0, err diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go index b6bcc7107..25693d711 100644 --- a/app/multitenant/collector.go +++ b/app/multitenant/collector.go @@ -3,6 +3,7 @@ package multitenant // Collect reports from probes per-tenant, and supply them to queriers on demand import ( + "encoding/json" "fmt" "io" "net" @@ -46,6 +47,48 @@ func (c *awsCollector) isCollector() bool { return c.cfg.StoreInterval != 0 } +func (c *awsCollector) hasReportsFromLive(ctx context.Context, userid string) (bool, error) { + span, ctx := opentracing.StartSpanFromContext(ctx, "hasReportsFromLive") + defer span.Finish() + if c.isCollector() { + e, found := c.pending.Load(userid) + if !found { + return false, nil + } + entry := e.(*pendingEntry) + entry.Lock() + defer entry.Unlock() + if entry.report != nil { + return true, nil + } + for _, v := range entry.older { + if v != nil { + return true, nil + } + } + return false, nil + } + // We are a querier: ask each collector if it has any + // (serially, since we will bail out on the first one that has reports) + addrs := resolve(c.cfg.CollectorAddr) + for _, addr := range addrs { + body, err := oneCall(ctx, addr, "/api/probes?sparse=true", userid) + if err != nil { + return false, err + } + var hasReports bool + decoder := json.NewDecoder(body) + if err := decoder.Decode(&hasReports); err != nil { + log.Errorf("Error encoding response: %v", err) + } + body.Close() + if hasReports { + return true, nil + } + } + return false, nil +} + func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]report.Report, error) { span, ctx := opentracing.StartSpanFromContext(ctx, "reportsFromLive") defer span.Finish() From 9b6202326685d8955c98df16204682fb97450f60 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Mon, 12 Apr 2021 19:47:07 +0000 Subject: [PATCH 09/11] Do REST calls from to collectors in parallel --- app/multitenant/collector.go | 43 ++++++++++++++++++++++++------------ 1 file changed, 29 insertions(+), 14 deletions(-) diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go index 25693d711..466dd3130 100644 --- a/app/multitenant/collector.go +++ b/app/multitenant/collector.go @@ -18,6 +18,7 @@ import ( log "github.com/sirupsen/logrus" "github.com/weaveworks/common/user" "github.com/weaveworks/scope/report" + "golang.org/x/sync/errgroup" ) // if StoreInterval is set, reports are merged into here and held until flushed to store @@ -115,22 +116,36 @@ func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]re // We are a querier: fetch the most up-to-date reports from collectors // TODO: resolve c.collectorAddress periodically instead of every time we make a call addrs := resolve(c.cfg.CollectorAddr) - ret := make([]report.Report, 0, len(addrs)) + reports := make([]*report.Report, len(addrs)) // make a call to each collector and fetch its data for this userid - // TODO: do them in parallel - for _, addr := range addrs { - body, err := oneCall(ctx, addr, "/api/report", userid) - if err != nil { - log.Warnf("error calling '%s': %v", addr, err) - continue + g, ctx := errgroup.WithContext(ctx) + for i, addr := range addrs { + i, addr := i, addr // https://golang.org/doc/faq#closures_and_goroutines + g.Go(func() error { + body, err := oneCall(ctx, addr, "/api/report", userid) + if err != nil { + log.Warnf("error calling '%s': %v", addr, err) + return nil + } + reports[i], err = report.MakeFromBinary(ctx, body, false, true) + body.Close() + if err != nil { + log.Warnf("error decoding: %v", err) + return nil + } + return nil + }) + } + if err := g.Wait(); err != nil { + return nil, err + } + + // dereference pointers into the expected return format + ret := make([]report.Report, 0, len(addrs)) + for _, rpt := range reports { + if rpt != nil { + ret = append(ret, *rpt) } - rpt, err := report.MakeFromBinary(ctx, body, false, true) - body.Close() - if err != nil { - log.Warnf("error decoding: %v", err) - continue - } - ret = append(ret, *rpt) } return ret, nil From 5856f372db558d1281d05111bbad2bcdf6905f3e Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 18 Apr 2021 17:33:18 +0000 Subject: [PATCH 10/11] multitenant: resolve collectors less frequently DNS records don't change that fast --- app/multitenant/aws_collector.go | 3 +++ app/multitenant/collector.go | 13 ++++++++----- 2 files changed, 11 insertions(+), 5 deletions(-) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 0a387e5c1..c18966125 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -156,6 +156,9 @@ type awsCollector struct { nats *nats.Conn waitersLock sync.Mutex waiters map[watchKey]*nats.Subscription + + collectors []string + lastResolved time.Time } // Shortcut reports: diff --git a/app/multitenant/collector.go b/app/multitenant/collector.go index 466dd3130..3c23373d6 100644 --- a/app/multitenant/collector.go +++ b/app/multitenant/collector.go @@ -10,6 +10,7 @@ import ( "net/http" "strconv" "sync" + "time" "context" @@ -114,12 +115,14 @@ func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]re } // We are a querier: fetch the most up-to-date reports from collectors - // TODO: resolve c.collectorAddress periodically instead of every time we make a call - addrs := resolve(c.cfg.CollectorAddr) - reports := make([]*report.Report, len(addrs)) + if time.Since(c.lastResolved) > time.Second*5 { + c.collectors = resolve(c.cfg.CollectorAddr) + c.lastResolved = time.Now() + } + reports := make([]*report.Report, len(c.collectors)) // make a call to each collector and fetch its data for this userid g, ctx := errgroup.WithContext(ctx) - for i, addr := range addrs { + for i, addr := range c.collectors { i, addr := i, addr // https://golang.org/doc/faq#closures_and_goroutines g.Go(func() error { body, err := oneCall(ctx, addr, "/api/report", userid) @@ -141,7 +144,7 @@ func (c *awsCollector) reportsFromLive(ctx context.Context, userid string) ([]re } // dereference pointers into the expected return format - ret := make([]report.Report, 0, len(addrs)) + ret := make([]report.Report, 0, len(reports)) for _, rpt := range reports { if rpt != nil { ret = append(ret, *rpt) From ced99f5008aa3f65661e99756e56f649f6eb1954 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Sun, 18 Apr 2021 18:57:12 +0000 Subject: [PATCH 11/11] multitenant: serialise report to buffer before sending Seems to be faster --- app/server_helpers.go | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/app/server_helpers.go b/app/server_helpers.go index 42d19fdfb..90f03d9d7 100644 --- a/app/server_helpers.go +++ b/app/server_helpers.go @@ -1,12 +1,14 @@ package app import ( + "bytes" "context" "net/http" "strings" opentracing "github.com/opentracing/opentracing-go" "github.com/ugorji/go/codec" + "github.com/weaveworks/scope/report" log "github.com/sirupsen/logrus" ) @@ -37,15 +39,20 @@ func respondWith(ctx context.Context, w http.ResponseWriter, code int, response // Similar to the above function, but respect the request's Accept header. // Possibly we should do a complete parse of Accept, but for now just rudimentary check -func respondWithReport(ctx context.Context, w http.ResponseWriter, req *http.Request, response interface{}) { +func respondWithReport(ctx context.Context, w http.ResponseWriter, req *http.Request, response report.Report) { accept := req.Header.Get("Accept") if strings.HasPrefix(accept, "application/msgpack") { - w.Header().Set("Content-Type", "application/msgpack") - w.WriteHeader(http.StatusOK) - encoder := codec.NewEncoder(w, &codec.MsgpackHandle{}) + buf := bytes.Buffer{} + encoder := codec.NewEncoder(&buf, &codec.MsgpackHandle{}) if err := encoder.Encode(response); err != nil { log.Errorf("Error encoding response: %v", err) } + if span := opentracing.SpanFromContext(ctx); span != nil { + span.LogKV("encoded-size", len(buf.Bytes())) + } + w.Header().Set("Content-Type", "application/msgpack") + w.WriteHeader(http.StatusOK) + w.Write(buf.Bytes()) return }