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