diff --git a/app/collector.go b/app/collector.go index 18b67bcc8..e7a81f08a 100644 --- a/app/collector.go +++ b/app/collector.go @@ -24,9 +24,13 @@ type Reporter interface { } // Adder is something that can accept reports. It's a convenient interface for -// parts of the app, and several experimental components. +// parts of the app, and several experimental components. It takes the following +// arguments: +// - context.Context: the request context +// - report.Report: the deserialised report +// - []byte: the serialised report (as gzip'd msgpack) type Adder interface { - Add(context.Context, report.Report) error + Add(context.Context, report.Report, []byte) error } // A Collector is a Reporter and an Adder @@ -88,7 +92,7 @@ func NewCollector(window time.Duration) Collector { } // Add adds a report to the collector's internal state. It implements Adder. -func (c *collector) Add(_ context.Context, rpt report.Report) error { +func (c *collector) Add(_ context.Context, rpt report.Report, _ []byte) error { c.mtx.Lock() defer c.mtx.Unlock() c.reports = append(c.reports, rpt) @@ -147,7 +151,7 @@ type StaticCollector report.Report func (c StaticCollector) Report(context.Context) (report.Report, error) { return report.Report(c), nil } // Add adds a report to the collector's internal state. It implements Adder. -func (c StaticCollector) Add(context.Context, report.Report) error { return nil } +func (c StaticCollector) Add(context.Context, report.Report, []byte) error { return nil } // WaitOn lets other components wait on a new report being received. It // implements Reporter. diff --git a/app/collector_test.go b/app/collector_test.go index 52566ea43..bfcfcfde2 100644 --- a/app/collector_test.go +++ b/app/collector_test.go @@ -32,7 +32,7 @@ func TestCollector(t *testing.T) { t.Error(test.Diff(want, have)) } - c.Add(ctx, r1) + c.Add(ctx, r1, nil) have, err = c.Report(ctx) if err != nil { t.Error(err) @@ -41,7 +41,7 @@ func TestCollector(t *testing.T) { t.Error(test.Diff(want, have)) } - c.Add(ctx, r2) + c.Add(ctx, r2, nil) merged := report.MakeReport() merged = merged.Merge(r1) merged = merged.Merge(r2) @@ -75,7 +75,7 @@ func TestCollectorExpire(t *testing.T) { // Now check an added report is returned r1 := report.MakeReport() r1.Endpoint.AddNode(report.MakeNode("foo")) - c.Add(ctx, r1) + c.Add(ctx, r1, nil) have, err = c.Report(ctx) if err != nil { t.Error(err) diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index 67e2f0ef5..2a32b82c7 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -354,7 +354,7 @@ func calculateReportKey(rowKey, colKey string) (string, error) { return fmt.Sprintf("%x/%s", rowKeyHash.Sum(nil), colKey), nil } -func (c *awsCollector) Add(ctx context.Context, rep report.Report) error { +func (c *awsCollector) Add(ctx context.Context, rep report.Report, buf []byte) error { userid, err := c.userIDer(ctx) if err != nil { return err @@ -367,7 +367,7 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report) error { return err } - reportSize, err := c.s3.StoreReport(reportKey, &rep) + reportSize, err := c.s3.StoreReportBytes(reportKey, buf) if err != nil { return err } @@ -375,7 +375,7 @@ func (c *awsCollector) Add(ctx context.Context, rep report.Report) error { // third, put it in memcache if c.memcache != nil { - _, err = c.memcache.StoreReport(reportKey, &rep) + _, err = c.memcache.StoreReportBytes(reportKey, buf) if err != nil { // NOTE: We don't abort here because failing to store in memcache // doesn't actually break anything else -- it's just an diff --git a/app/multitenant/memcache_client.go b/app/multitenant/memcache_client.go index 0a42b9953..9d3e165cc 100644 --- a/app/multitenant/memcache_client.go +++ b/app/multitenant/memcache_client.go @@ -205,13 +205,11 @@ func (c *MemcacheClient) FetchReports(keys []string) (map[string]report.Report, return reports, missing, nil } -// StoreReport serializes and stores a report. -func (c *MemcacheClient) StoreReport(key string, report *report.Report) (int, error) { - var buf bytes.Buffer - report.WriteBinary(&buf, c.compressionLevel) +// StoreReportBytes stores a report. +func (c *MemcacheClient) StoreReportBytes(key string, rpt []byte) (int, error) { err := instrument.TimeRequestHistogramStatus("Put", memcacheRequestDuration, memcacheStatusCode, func() error { - item := memcache.Item{Key: key, Value: buf.Bytes(), Expiration: c.expiration} + item := memcache.Item{Key: key, Value: rpt, Expiration: c.expiration} return c.client.Set(&item) }) - return buf.Len(), err + return len(rpt), err } diff --git a/app/multitenant/s3_client.go b/app/multitenant/s3_client.go index 9e1aaaeac..99726b196 100644 --- a/app/multitenant/s3_client.go +++ b/app/multitenant/s3_client.go @@ -2,7 +2,6 @@ package multitenant import ( "bytes" - "compress/gzip" "github.com/aws/aws-sdk-go/aws" "github.com/aws/aws-sdk-go/aws/session" @@ -85,19 +84,15 @@ func (store *S3Store) fetchReport(key string) (*report.Report, error) { return report.MakeFromBinary(resp.Body) } -// StoreReport serializes and stores a report. -// -// Returns the size of the report. This only equals bytes written if err is nil. -func (store *S3Store) StoreReport(key string, report *report.Report) (int, error) { - var buf bytes.Buffer - report.WriteBinary(&buf, gzip.BestCompression) +// StoreReportBytes stores a report. +func (store *S3Store) StoreReportBytes(key string, buf []byte) (int, error) { err := instrument.TimeRequestHistogram("Put", s3RequestDuration, func() error { _, err := store.s3.PutObject(&s3.PutObjectInput{ - Body: bytes.NewReader(buf.Bytes()), + Body: bytes.NewReader(buf), Bucket: aws.String(store.bucketName), Key: aws.String(key), }) return err }) - return buf.Len(), err + return len(buf), err } diff --git a/app/router.go b/app/router.go index 1e200533c..d468f3aa5 100644 --- a/app/router.go +++ b/app/router.go @@ -1,7 +1,10 @@ package app import ( + "bytes" + "compress/gzip" "fmt" + "io" "net/http" "net/url" "strings" @@ -112,16 +115,22 @@ func RegisterReportPostHandler(a Adder, router *mux.Router) { post.HandleFunc("/api/report", requestContextDecorator(func(ctx context.Context, w http.ResponseWriter, r *http.Request) { var ( rpt report.Report - reader = r.Body + buf bytes.Buffer + reader = io.TeeReader(r.Body, &buf) ) gzipped := strings.Contains(r.Header.Get("Content-Encoding"), "gzip") + if !gzipped { + reader = io.TeeReader(r.Body, gzip.NewWriter(&buf)) + } + contentType := r.Header.Get("Content-Type") + isMsgpack := strings.HasPrefix(contentType, "application/msgpack") var handle codec.Handle switch { case strings.HasPrefix(contentType, "application/json"): handle = &codec.JsonHandle{} - case strings.HasPrefix(contentType, "application/msgpack"): + case isMsgpack: handle = &codec.MsgpackHandle{} default: respondWith(w, http.StatusBadRequest, fmt.Errorf("Unsupported Content-Type: %v", contentType)) @@ -133,7 +142,13 @@ func RegisterReportPostHandler(a Adder, router *mux.Router) { return } - if err := a.Add(ctx, rpt); err != nil { + // a.Add(..., buf) assumes buf is gzip'd msgpack + if !isMsgpack { + buf = bytes.Buffer{} + rpt.WriteBinary(&buf, gzip.BestCompression) + } + + if err := a.Add(ctx, rpt, buf.Bytes()); err != nil { log.Errorf("Error Adding report: %v", err) respondWith(w, http.StatusInternalServerError, err) return diff --git a/docker/Dockerfile b/docker/Dockerfile index 997294ebd..029177008 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -2,7 +2,7 @@ FROM alpine:3.3 MAINTAINER Weaveworks Inc LABEL works.weave.role=system WORKDIR /home/weave -RUN echo "http://dl-3.alpinelinux.org/alpine/edge/testing" >>/etc/apk/repositories && \ +RUN echo "http://dl-cdn.alpinelinux.org/alpine/edge/community" >>/etc/apk/repositories && \ apk add --update bash runit conntrack-tools iproute2 util-linux curl && \ rm -rf /var/cache/apk/* ADD ./docker.tgz /