Don't reencode reports in the collector (#1819)

* Don't reencode reports in the collector

* Review feedback

* Fix comment

* Update alpine URLs so it will build

* Fix tests
This commit is contained in:
Tom Wilkie
2016-08-22 17:37:41 +01:00
committed by GitHub
parent 4cb002e360
commit 32edfc9112
7 changed files with 41 additions and 29 deletions

View File

@@ -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.

View File

@@ -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)

View File

@@ -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

View File

@@ -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
}

View File

@@ -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
}

View File

@@ -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

View File

@@ -2,7 +2,7 @@ FROM alpine:3.3
MAINTAINER Weaveworks Inc <help@weave.works>
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 /