diff --git a/storage/aws/aws.go b/storage/aws/aws.go index a448275ab..2ee96a26a 100644 --- a/storage/aws/aws.go +++ b/storage/aws/aws.go @@ -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), @@ -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 } diff --git a/storage/aws/aws_test.go b/storage/aws/aws_test.go index bfa35032d..9ec511627 100644 --- a/storage/aws/aws_test.go +++ b/storage/aws/aws_test.go @@ -29,6 +29,8 @@ import ( "errors" "flag" "fmt" + "net/http" + "net/http/httptest" "os" "reflect" "strings" @@ -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" @@ -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 /, +// not //. 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, `PreconditionFailedAt least one of the pre-conditions you specified did not hold`) + 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, `NoSuchKeyThe specified key does not exist.`) + 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) + } +} diff --git a/storage/gcp/gcp.go b/storage/gcp/gcp.go index 13fea6d11..8f95ae1f7 100644 --- a/storage/gcp/gcp.go +++ b/storage/gcp/gcp.go @@ -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 { @@ -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 { @@ -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 }) diff --git a/storage/gcp/gcp_test.go b/storage/gcp/gcp_test.go index 33540a268..1e4cc91d4 100644 --- a/storage/gcp/gcp_test.go +++ b/storage/gcp/gcp_test.go @@ -22,6 +22,9 @@ import ( "fmt" "log" "log/slog" + "net/http" + "net/http/httptest" + "net/url" "os" "reflect" "strings" @@ -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 /, not //. +// 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//o/?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 {