Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 53 additions & 10 deletions storage/gcp/gcp.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ package gcp

import (
"bytes"
"compress/gzip"
"context"
"encoding/gob"
"errors"
Expand Down Expand Up @@ -83,6 +84,7 @@ const (
ckptContType = "text/plain; charset=utf-8"
logCacheControl = "max-age=604800,immutable"
ckptCacheControl = "no-cache"
bundleContEnc = "gzip"

DefaultIntegrationSizeLimit = 5 * 4096

Expand Down Expand Up @@ -519,7 +521,7 @@ func (a *Appender) updateCheckpoint(ctx context.Context, size uint64, root []byt
// objStore describes a type which can store and retrieve objects.
type objStore interface {
getObject(ctx context.Context, obj string) ([]byte, *gcs.ReaderObjectAttrs, error)
setObject(ctx context.Context, obj string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string) error
setObject(ctx context.Context, obj string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string, contEnc string) error
deleteObjectsWithPrefix(ctx context.Context, prefix string) error
}

Expand All @@ -530,7 +532,7 @@ type logResourceStore struct {
}

func (lrs *logResourceStore) setCheckpoint(ctx context.Context, cpRaw []byte) error {
return lrs.objStore.setObject(ctx, layout.CheckpointPath, cpRaw, nil, ckptContType, ckptCacheControl)
return lrs.objStore.setObject(ctx, layout.CheckpointPath, cpRaw, nil, ckptContType, ckptCacheControl, "")
}

func (lrs *logResourceStore) getCheckpoint(ctx context.Context) ([]byte, error) {
Expand All @@ -550,7 +552,7 @@ func (s *logResourceStore) setTile(ctx context.Context, level, index uint64, par
start := time.Now()

tPath := layout.TilePath(level, index, partial)
err := s.objStore.setObject(ctx, tPath, data, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl)
err := s.objStore.setObject(ctx, tPath, data, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl, "")
opsHistogram.Record(ctx, time.Since(start).Milliseconds(), metric.WithAttributes(opNameKey.String("writeTile")))
return err
}
Expand Down Expand Up @@ -628,7 +630,7 @@ func (s *logResourceStore) setEntryBundle(ctx context.Context, bundleIndex uint6
// Note that setObject does an idempotent interpretation of DoesNotExist - it only
// returns an error if the named object exists _and_ contains different data to what's
// passed in here.
if err := s.objStore.setObject(ctx, objName, bundleRaw, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl); err != nil {
if err := s.objStore.setObject(ctx, objName, bundleRaw, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl, bundleContEnc); err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Might be worth popping a comment explaining that if a client doesn't set Accept-Encoding, then GCS (I think?) will decompress data.

return fmt.Errorf("setObject(%q): %v", objName, err)

}
Expand Down Expand Up @@ -1338,9 +1340,33 @@ func (s *gcsStorage) getObject(ctx context.Context, obj string) ([]byte, *gcs.Re
return nil, nil, fmt.Errorf("getObject: failed to create reader for object %q in bucket %q: %w", obj, s.bucket, err)
}

d, err := io.ReadAll(r)
if err != nil {
return nil, nil, fmt.Errorf("failed to read %q: %v", obj, err)
// The HTTP GCS client transparently decompresses Content-Encoding: gzip objects
// (setting r.Attrs.Decompressed = true and clearing r.Attrs.ContentEncoding),
// whereas the gRPC GCS client returns the raw compressed bytes with
// r.Attrs.ContentEncoding == "gzip" and r.Attrs.Decompressed == false.
var d []byte
if r.Attrs.ContentEncoding == "gzip" && !r.Attrs.Decompressed {
zr, err := gzip.NewReader(r)
if err != nil {
_ = r.Close()
return nil, nil, fmt.Errorf("getObject: failed to create gzip reader for %q: %v", obj, err)
}
d, err = io.ReadAll(zr)
if err != nil {
_ = zr.Close()
_ = r.Close()
return nil, nil, fmt.Errorf("failed to read %q: %v", obj, err)
}
if err := zr.Close(); err != nil {
_ = r.Close()
return nil, nil, fmt.Errorf("failed to close gzip reader for %q: %v", obj, err)
}
} else {
d, err = io.ReadAll(r)
if err != nil {
_ = r.Close()
return nil, nil, fmt.Errorf("failed to read %q: %v", obj, err)
}
}
return d, &r.Attrs, r.Close()
})
Expand All @@ -1354,7 +1380,7 @@ func (s *gcsStorage) getObject(ctx context.Context, obj string) ([]byte, *gcs.Re
// Note that when preconditions are specified and are not met, an error will be returned *unless*
// the currently stored data is bit-for-bit identical to the data to-be-written.
// This is intended to provide idempotentency for writes.
func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string) error {
func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string, contEnc string) error {
return otel.TraceErr(ctx, "tessera.storage.gcp.setObject", tracer, func(ctx context.Context, span trace.Span) error {
// objName deliberately stays unprefixed so the precondition-failed
// recovery getObject call below, which applies bucketPrefix itself,
Expand All @@ -1366,6 +1392,22 @@ func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte,

span.SetAttributes(objectPathKey.String(prefixedName))

writeData := data
if contEnc != "" {
if contEnc != "gzip" {
return fmt.Errorf("unsupported content encoding %q for %q", contEnc, prefixedName)
}
var b bytes.Buffer
zw := gzip.NewWriter(&b)
if _, err := zw.Write(data); err != nil {
return fmt.Errorf("failed to gzip data for %q: %v", prefixedName, err)
}
if err := zw.Close(); err != nil {
return fmt.Errorf("failed to close gzip writer for %q: %v", prefixedName, err)
}
writeData = b.Bytes()
}

bkt := s.gcsClient.Bucket(s.bucket)
obj := bkt.Object(prefixedName)

Expand All @@ -1378,9 +1420,10 @@ func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte,
}
w.ContentType = contType
w.CacheControl = cacheCtl
w.ContentEncoding = contEnc
// Limit the amount of memory used for buffers, see https://pkg.go.dev/cloud.google.com/go/storage#Writer
w.ChunkSize = len(data) + 1024
if _, err := w.Write(data); err != nil {
w.ChunkSize = len(writeData) + 1024
if _, err := w.Write(writeData); err != nil {
return fmt.Errorf("failed to write object %q to bucket %q: %w", prefixedName, s.bucket, err)
}

Expand Down
Loading
Loading