Skip to content
Merged
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
18 changes: 11 additions & 7 deletions storage/aws/aws.go
Original file line number Diff line number Diff line change
Expand Up @@ -1531,13 +1531,17 @@ func (s *s3Storage) setObject(ctx context.Context, objName string, data []byte,
// 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 *s3Storage) setObjectIfNoneMatch(ctx context.Context, objName string, data []byte, contType string, cacheControl string) error {
// objName deliberately stays unprefixed so the precondition-failed
// recovery getObject call below, which applies bucketPrefix itself,
// can use it directly.
prefixedName := objName
if s.bucketPrefix != "" {
objName = filepath.Join(s.bucketPrefix, objName)
prefixedName = filepath.Join(s.bucketPrefix, objName)
}

put := &s3.PutObjectInput{
Bucket: aws.String(s.bucket),
Key: aws.String(objName),
Key: aws.String(prefixedName),
Body: bytes.NewReader(data),
ContentType: aws.String(contType),
CacheControl: aws.String(cacheControl),
Expand All @@ -1554,18 +1558,18 @@ func (s *s3Storage) setObjectIfNoneMatch(ctx context.Context, objName string, da
if errors.As(err, &apiErr) && apiErr.ErrorCode() == "PreconditionFailed" {
existing, err := s.getObject(ctx, objName)
if err != nil {
return fmt.Errorf("failed to fetch existing content for %q: %v", objName, err)
return fmt.Errorf("failed to fetch existing content for %q: %v", prefixedName, err)
}
if !bytes.Equal(existing, data) {
slog.ErrorContext(ctx, "Resource non-idempotent writen", slog.String("objname", objName), slog.String("diff", cmp.Diff(existing, data)))
return fmt.Errorf("precondition failed: resource content for %q differs from data to-be-written", objName)
slog.ErrorContext(ctx, "Resource non-idempotent writen", slog.String("objname", prefixedName), slog.String("diff", cmp.Diff(existing, data)))
return fmt.Errorf("precondition failed: resource content for %q differs from data to-be-written", prefixedName)
}

slog.DebugContext(ctx, "setObjectIfNoneMatch: identical resource already exists. Continuing", slog.String("objname", objName))
slog.DebugContext(ctx, "setObjectIfNoneMatch: identical resource already exists. Continuing", slog.String("objname", prefixedName))
return nil
}

return fmt.Errorf("failed to write object %q to bucket %q: %w", objName, s.bucket, err)
return fmt.Errorf("failed to write object %q to bucket %q: %w", prefixedName, s.bucket, err)
}
return nil
}
Expand Down
79 changes: 79 additions & 0 deletions storage/aws/aws_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,8 @@ import (
"errors"
"flag"
"fmt"
"net/http"
"net/http/httptest"
"os"
"reflect"
"strings"
Expand All @@ -38,6 +40,9 @@ import (

"log/slog"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/aws/smithy-go"
"github.com/google/go-cmp/cmp"
Expand Down Expand Up @@ -833,3 +838,77 @@ func defaultMerkleLeafHasher(bundle []byte) ([][]byte, error) {
}
return r, nil
}

// TestSetObjectIfNoneMatchBucketPrefixIdempotentRecovery exercises
// s3Storage.setObjectIfNoneMatch's precondition-failed recovery path with a
// non-empty BucketPrefix: the recovery getObject must read <prefix>/<name>,
// not <prefix>/<prefix>/<name>. setObjectIfNoneMatch and getObject each apply
// bucketPrefix once, so setObjectIfNoneMatch must hand getObject the
// unprefixed name. The fake server returns 412 on the write and serves
// identical bytes at the single-prefixed key, so setObjectIfNoneMatch should
// return nil (idempotent success).
func TestSetObjectIfNoneMatchBucketPrefixIdempotentRecovery(t *testing.T) {
const (
bucket = "test-bucket"
prefix = "some/prefix"
objName = "tile/0/000"
)
data := []byte("tile-bytes")
wantKey := prefix + "/" + objName

// This fake pins the aws-sdk-go-v2 S3 client's REST-XML wire shape
// (path-style requests, XML error bodies). That coupling is deliberate:
// the bug under test sits below the objStore interface, so a real
// *s3.Client is the only way to exercise it.
var mu sync.Mutex
var gotReadKey string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
key := strings.TrimPrefix(r.URL.Path, "/"+bucket+"/")
switch r.Method {
case http.MethodPut:
if key != wantKey {
t.Errorf("conditional write PUT key %q, want %q", key, wantKey)
}
// Conditional write: object already exists, precondition fails.
w.Header().Set("Content-Type", "application/xml")
w.WriteHeader(http.StatusPreconditionFailed)
_, _ = fmt.Fprint(w, `<?xml version="1.0" encoding="UTF-8"?><Error><Code>PreconditionFailed</Code><Message>At least one of the pre-conditions you specified did not hold</Message></Error>`)
case http.MethodGet:
mu.Lock()
gotReadKey = key
mu.Unlock()
if key != wantKey {
w.Header().Set("Content-Type", "application/xml")
w.WriteHeader(http.StatusNotFound)
_, _ = fmt.Fprint(w, `<?xml version="1.0" encoding="UTF-8"?><Error><Code>NoSuchKey</Code><Message>The specified key does not exist.</Message></Error>`)
return
}
w.Header().Set("Content-Type", "application/octet-stream")
w.WriteHeader(http.StatusOK)
_, _ = w.Write(data)
default:
t.Errorf("unexpected request: %s %s", r.Method, r.URL)
http.Error(w, "unexpected", http.StatusNotImplemented)
}
}))
defer srv.Close()

c := s3.New(s3.Options{
BaseEndpoint: aws.String(srv.URL),
UsePathStyle: true,
Region: "us-east-1",
Credentials: credentials.NewStaticCredentialsProvider("test", "test", ""),
})

s := &s3Storage{s3Client: c, bucket: bucket, bucketPrefix: prefix}
err := s.setObjectIfNoneMatch(context.Background(), objName, data, "application/octet-stream", "")
mu.Lock()
got := gotReadKey
mu.Unlock()
if err != nil {
t.Fatalf("setObjectIfNoneMatch: want idempotent success (nil), got %v (recovery read was for %q)", err, got)
}
if got != wantKey {
t.Fatalf("recovery getObject read %q, want %q (double-prefix bug)", got, wantKey)
}
}
22 changes: 13 additions & 9 deletions storage/gcp/gcp.go
Original file line number Diff line number Diff line change
Expand Up @@ -1355,14 +1355,18 @@ func (s *gcsStorage) getObject(ctx context.Context, obj string) ([]byte, *gcs.Re
// 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 {
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,
// can use it directly.
prefixedName := objName
if s.bucketPrefix != "" {
objName = filepath.Join(s.bucketPrefix, objName)
prefixedName = filepath.Join(s.bucketPrefix, objName)
}

span.SetAttributes(objectPathKey.String(objName))
span.SetAttributes(objectPathKey.String(prefixedName))

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

var w *gcs.Writer
if cond == nil {
Expand All @@ -1376,7 +1380,7 @@ func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte,
// 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 {
return fmt.Errorf("failed to write object %q to bucket %q: %w", objName, s.bucket, err)
return fmt.Errorf("failed to write object %q to bucket %q: %w", prefixedName, s.bucket, err)
}

if err := w.Close(); err != nil {
Expand All @@ -1395,20 +1399,20 @@ func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte,
if preconditionFailed {
existing, existingAttr, err := s.getObject(ctx, objName)
if err != nil {
return fmt.Errorf("failed to fetch existing content for %q: %v", objName, err)
return fmt.Errorf("failed to fetch existing content for %q: %v", prefixedName, err)
}
if !bytes.Equal(existing, data) {
span.AddEvent("Non-idempotent write")
slog.ErrorContext(ctx, "Resource non-idempotent write", slog.String("objName", objName), slog.String("diff", cmp.Diff(existing, data)))
return fmt.Errorf("precondition failed: resource content for %q (@%d) differs from data to-be-written", objName, existingAttr.Generation)
slog.ErrorContext(ctx, "Resource non-idempotent write", slog.String("objName", prefixedName), slog.String("diff", cmp.Diff(existing, data)))
return fmt.Errorf("precondition failed: resource content for %q (@%d) differs from data to-be-written", prefixedName, existingAttr.Generation)
}

span.AddEvent("Idempotent write")
slog.DebugContext(ctx, "setObject: identical resource already exists", slog.String("objName", objName))
slog.DebugContext(ctx, "setObject: identical resource already exists", slog.String("objName", prefixedName))
return nil
}

return fmt.Errorf("failed to close write on %q: %v", objName, err)
return fmt.Errorf("failed to close write on %q: %v", prefixedName, err)
}
return nil
})
Expand Down
83 changes: 83 additions & 0 deletions storage/gcp/gcp_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ import (
"fmt"
"log"
"log/slog"
"net/http"
"net/http/httptest"
"net/url"
"os"
"reflect"
"strings"
Expand Down Expand Up @@ -941,6 +944,86 @@ func (m *memObjStore) deleteObjectsWithPrefix(_ context.Context, prefix string)
return nil
}

// TestSetObjectBucketPrefixIdempotentRecovery exercises gcsStorage.setObject's
// precondition-failed recovery path with a non-empty BucketPrefix: the
// recovery getObject must read <prefix>/<name>, not <prefix>/<prefix>/<name>.
// setObject and getObject each apply bucketPrefix once, so setObject must
// hand getObject the unprefixed name. The fake server returns 412 on the
// write and serves identical bytes at the single-prefixed path, so setObject
// should return nil (idempotent success).
func TestSetObjectBucketPrefixIdempotentRecovery(t *testing.T) {
const (
bucket = "test-bucket"
prefix = "some/prefix"
objName = "tile/0/000"
)
data := []byte("tile-bytes")
wantObj := prefix + "/" + objName

// This fake pins the cloud.google.com/go/storage client's HTTP JSON-API
// wire shape (multipart upload, 412 parsing, JSON media-read path). That
// coupling is deliberate: the bug under test sits below the objStore
// interface, so a real *gcs.Client is the only way to exercise it; the
// cost is possible breakage on a client-library upgrade that reshapes
// these requests.
var mu sync.Mutex
var gotReadObj string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodPost && strings.HasPrefix(r.URL.Path, "/upload/"):
// Conditional write: object already exists, precondition fails.
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusPreconditionFailed)
_, _ = fmt.Fprint(w, `{"error":{"code":412,"message":"conditionNotMet"}}`)
case r.Method == http.MethodGet:
// JSON-API media read. Tolerate both path shapes the client
// library has been observed to emit against an emulator endpoint:
// [/download]/storage/v1/b/<bucket>/o/<url-encoded object>?alt=media
p := strings.TrimPrefix(r.URL.Path, "/download")
got, _ := url.PathUnescape(strings.TrimPrefix(p, "/storage/v1/b/"+bucket+"/o/"))
mu.Lock()
gotReadObj = got
mu.Unlock()
if got != wantObj {
http.Error(w, "object not found", http.StatusNotFound)
return
}
w.Header().Set("Content-Type", "application/octet-stream")
w.Header().Set("X-Goog-Generation", "1")
w.WriteHeader(http.StatusOK)
_, _ = w.Write(data)
default:
t.Errorf("unexpected request: %s %s", r.Method, r.URL)
http.Error(w, "unexpected", http.StatusNotImplemented)
}
}))
defer srv.Close()

t.Setenv("STORAGE_EMULATOR_HOST", strings.TrimPrefix(srv.URL, "http://"))
ctx := context.Background()
c, err := gcs.NewClient(ctx, gcs.WithJSONReads())
if err != nil {
t.Fatalf("NewClient: %v", err)
}
defer func() {
if err := c.Close(); err != nil {
t.Logf("Close: %v", err)
}
}()

s := &gcsStorage{gcsClient: c, bucket: bucket, bucketPrefix: prefix}
err = s.setObject(ctx, objName, data, &gcs.Conditions{DoesNotExist: true}, "application/octet-stream", "")
mu.Lock()
got := gotReadObj
mu.Unlock()
if err != nil {
t.Fatalf("setObject: want idempotent success (nil), got %v (recovery read was for %q)", err, got)
}
if got != wantObj {
t.Fatalf("recovery getObject read %q, want %q (double-prefix bug)", got, wantObj)
}
}

func mustGenerateKeys(t *testing.T) (note.Signer, note.Verifier) {
sk, vk, err := note.GenerateKey(nil, "testlog")
if err != nil {
Expand Down
Loading