From 13269e81105fd4d18d298b6d10698e23ecade921 Mon Sep 17 00:00:00 2001 From: Jonathan Lange Date: Thu, 16 Jun 2016 18:17:27 +0100 Subject: [PATCH] Helper for reading & writing from binary --- app/multitenant/dynamo_collector.go | 21 ++------------- probe/appclient/report_publisher.go | 10 +------ report/marshal.go | 42 +++++++++++++++++++++++++++++ 3 files changed, 45 insertions(+), 28 deletions(-) create mode 100644 report/marshal.go diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index c8cec9168..0d5181d14 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -2,7 +2,6 @@ package multitenant import ( "bytes" - "compress/gzip" "crypto/md5" "fmt" "io" @@ -18,7 +17,6 @@ import ( "github.com/bluele/gcache" "github.com/nats-io/nats" "github.com/prometheus/client_golang/prometheus" - "github.com/ugorji/go/codec" "golang.org/x/net/context" "github.com/weaveworks/scope/app" @@ -312,15 +310,7 @@ func (c *dynamoDBCollector) getNonCachedReport(reportKey string) (*report.Report if err != nil { return nil, err } - reader, err := gzip.NewReader(resp.Body) - if err != nil { - return nil, err - } - rep := report.MakeReport() - if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { - return nil, err - } - return &rep, nil + return report.MakeFromBinary(resp.Body) } func (c *dynamoDBCollector) getReports(userid string, row int64, start, end time.Time) ([]report.Report, error) { @@ -387,14 +377,7 @@ func (c *dynamoDBCollector) Add(ctx context.Context, rep report.Report) error { // first, encode the report into a buffer and record its size var buf bytes.Buffer - writer, err := gzip.NewWriterLevel(&buf, gzip.BestCompression) - if err != nil { - return err - } - if err := codec.NewEncoder(writer, &codec.MsgpackHandle{}).Encode(&rep); err != nil { - return err - } - writer.Close() + rep.WriteBinary(&buf) reportSize.Add(float64(buf.Len())) // second, put the report on s3 diff --git a/probe/appclient/report_publisher.go b/probe/appclient/report_publisher.go index 5e713ebb8..2e2134147 100644 --- a/probe/appclient/report_publisher.go +++ b/probe/appclient/report_publisher.go @@ -2,9 +2,6 @@ package appclient import ( "bytes" - "compress/gzip" - "github.com/ugorji/go/codec" - "github.com/weaveworks/scope/report" ) @@ -24,11 +21,6 @@ func NewReportPublisher(publisher Publisher) *ReportPublisher { // Publish serialises and compresses a report, then passes it to a publisher func (p *ReportPublisher) Publish(r report.Report) error { buf := &bytes.Buffer{} - gzwriter := gzip.NewWriter(buf) - if err := codec.NewEncoder(gzwriter, &codec.MsgpackHandle{}).Encode(r); err != nil { - return err - } - gzwriter.Close() // otherwise the content won't get flushed to the output stream - + r.WriteBinary(buf) return p.publisher.Publish(buf) } diff --git a/report/marshal.go b/report/marshal.go new file mode 100644 index 000000000..39fa7704d --- /dev/null +++ b/report/marshal.go @@ -0,0 +1,42 @@ +package report + +import ( + "compress/gzip" + "io" + + "github.com/ugorji/go/codec" +) + +// WriteBinary writes a Report as a gzipped msgpack. +func (rep Report) WriteBinary(w io.Writer) error { + gzwriter, err := gzip.NewWriterLevel(w, gzip.BestCompression) + if err != nil { + return err + } + if err = codec.NewEncoder(gzwriter, &codec.MsgpackHandle{}).Encode(&rep); err != nil { + return err + } + gzwriter.Close() // otherwise the content won't get flushed to the output stream + return nil +} + +// ReadBinary reads into a Report from a gzipped msgpack. +func (rep *Report) ReadBinary(r io.Reader) error { + reader, err := gzip.NewReader(r) + if err != nil { + return err + } + if err := codec.NewDecoder(reader, &codec.MsgpackHandle{}).Decode(&rep); err != nil { + return err + } + return nil +} + +// MakeFromBinary constructs a Report from a gzipped msgpack. +func MakeFromBinary(r io.Reader) (*Report, error) { + rep := MakeReport() + if err := rep.ReadBinary(r); err != nil { + return nil, err + } + return &rep, nil +}