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
2 changes: 1 addition & 1 deletion packages/shared/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ require (
github.com/RoaringBitmap/roaring/v2 v2.18.0
github.com/aws/aws-sdk-go-v2 v1.41.6
github.com/aws/aws-sdk-go-v2/config v1.32.12
github.com/aws/aws-sdk-go-v2/credentials v1.19.12
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.20.12
github.com/aws/aws-sdk-go-v2/service/ecr v1.44.0
github.com/aws/aws-sdk-go-v2/service/s3 v1.100.0
Expand Down Expand Up @@ -98,7 +99,6 @@ require (
github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect
github.com/armon/go-metrics v0.4.1 // indirect
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.9 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.19.12 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.20 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.22 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.22 // indirect
Expand Down
299 changes: 299 additions & 0 deletions packages/shared/pkg/storage/gcp_multipart_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package storage

import (
"bytes"
"crypto/md5"
"crypto/sha256"
"encoding/base64"
Expand All @@ -13,12 +14,15 @@ import (
"net/http/httptest"
"os"
"path/filepath"
"slices"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/hashicorp/go-retryablehttp"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -1194,3 +1198,298 @@ func TestGCPCompleteRetriesTruncatedBody(t *testing.T) {
require.Error(t, err, "a truncated complete response must fail, not commit")
require.Equal(t, int32(3), attempts.Load(), "a failed body read must be retried to the budget")
}

// gcsXMLBackend starts MinIO and opens its bucket for anonymous access so the
// GCS XML MultipartUploader (bearer-token auth, unverifiable by MinIO) can
// upload without request signing.
func gcsXMLBackend(t *testing.T) *s3TestBackend {
t.Helper()

backend := startMinioBackend(t)

policy := fmt.Sprintf(`{
"Version": "2012-10-17",
"Statement": [{
"Effect": "Allow",
"Principal": {"AWS": ["*"]},
"Action": ["s3:*"],
"Resource": ["arn:aws:s3:::%s", "arn:aws:s3:::%s/*"]
}]
}`, backend.bucket, backend.bucket)

_, err := backend.newClient(t, nil).PutBucketPolicy(t.Context(), &s3.PutBucketPolicyInput{
Bucket: aws.String(backend.bucket),
Policy: aws.String(policy),
})
require.NoError(t, err, "allow anonymous access on minio bucket")

return backend
}

// authStrippingTransport removes the Authorization header (the uploader's
// bearer token) so requests reach MinIO as anonymous — MinIO cannot validate
// Google OAuth tokens. Nothing else is adapted: in particular, MinIO's 411
// rejection of chunked uploads guards the uploader's explicit Content-Length
// (via multiSliceReader.Len) against regressions.
type authStrippingTransport struct {
inner http.RoundTripper
}

func (a authStrippingTransport) RoundTrip(req *http.Request) (*http.Response, error) {
clone := req.Clone(req.Context())
clone.Header.Del("Authorization")

return a.inner.RoundTrip(clone)
}

// gcsXMLUploader builds a MultipartUploader pointed at the MinIO backend,
// bypassing NewMultipartUploaderWithRetryConfig (which requires real Google
// credentials and hardcodes the production URL). transport is optional and
// sits between the retryable client and the network (e.g. fault injection).
func gcsXMLUploader(t *testing.T, backend *s3TestBackend, key string, metadata ObjectMetadata, transport http.RoundTripper) *MultipartUploader {
t.Helper()

if transport == nil {
transport = http.DefaultTransport
}

rc := retryablehttp.NewClient()
rc.RetryMax = 3
rc.RetryWaitMin = 10 * time.Millisecond
rc.RetryWaitMax = 50 * time.Millisecond
rc.Logger = nil
rc.HTTPClient = &http.Client{Transport: authStrippingTransport{inner: transport}}

return &MultipartUploader{
bucketName: backend.bucket,
objectName: key,
token: "test-token", // stripped by authStrippingTransport
client: rc,
retryConfig: DefaultRetryConfig(),
metadata: metadata,
baseURL: backend.endpoint + "/" + backend.bucket,
}
}

// anonymousGet fetches an object's raw bytes via unauthenticated HTTP.
func anonymousGet(t *testing.T, backend *s3TestBackend, key string) []byte {
t.Helper()

req, err := http.NewRequestWithContext(t.Context(), http.MethodGet, backend.endpoint+"/"+backend.bucket+"/"+key, nil)
require.NoError(t, err)

resp, err := http.DefaultClient.Do(req)
require.NoError(t, err)
defer resp.Body.Close()
require.Equal(t, http.StatusOK, resp.StatusCode)

body, err := io.ReadAll(resp.Body)
require.NoError(t, err)

return body
}
Comment thread
dobrac marked this conversation as resolved.

// TestGCSXMLPartUploaderContract drives MultipartUploader directly against
// MinIO's XML multipart implementation: out-of-order part numbers, a
// multi-slice part body, Content-MD5 validation by a real server, and
// ordered reassembly on Complete.
func TestGCSXMLPartUploaderContract(t *testing.T) {
t.Parallel()

backend := gcsXMLBackend(t)
key := testKey("gcs-part-contract")
up := gcsXMLUploader(t, backend, key, nil, nil)

// Part 1 (non-final) must be >= 5 MiB; two slices exercise the
// multi-slice hashing/streaming path. Part 2 (final) is tiny.
part1a := bytes.Repeat([]byte{0xC3}, 3*megabyte)
part1b := bytes.Repeat([]byte{0xD4}, 2*megabyte+512)
part2 := []byte("gcs-final-part")

require.NoError(t, up.Start(t.Context()))
// Upload out of order: final part first.
require.NoError(t, up.UploadPart(t.Context(), 2, part2))
require.NoError(t, up.UploadPart(t.Context(), 1, part1a, part1b))
require.NoError(t, up.Complete(t.Context()))
require.NoError(t, up.Close())

want := slices.Concat(part1a, part1b, part2)
got := anonymousGet(t, backend, key)
require.Equal(t, sha256.Sum256(want), sha256.Sum256(got),
"parts must reassemble in part-number order, not upload order")
}

// TestGCSXMLCompressedRoundTrip runs the production compressed upload path
// (storeFileCompressed) with the GCS XML uploader against MinIO, then
// verifies the stored blob decompresses back to the original.
func TestGCSXMLCompressedRoundTrip(t *testing.T) {
t.Parallel()

codecs := []struct {
codec CompressionType
level int
}{
{CompressionZstd, 2},
{CompressionLZ4, 0},
}

for _, tc := range codecs {
t.Run(tc.codec.String(), func(t *testing.T) {
t.Parallel()

const dataSize = 32 * megabyte
data := generateSemiRandomData(dataSize)
inputPath := writeTempFile(t, data)

backend := gcsXMLBackend(t)
key := testKey("gcs-compressed-" + tc.codec.String())

cfg := CompressConfig{
Enabled: true,
Type: tc.codec.String(),
Level: tc.level,
FrameSizeKB: 2 * 1024,
MinPartSizeMB: 5,
FrameEncodeWorkers: 4,
EncoderConcurrency: 1,
}

fullFT, checksum, err := storeFileCompressed(t.Context(), inputPath, cfg, 4, PutOptions{},
func(metadata ObjectMetadata) (partUploader, error) {
return gcsXMLUploader(t, backend, key, metadata, nil), nil
})
require.NoError(t, err)
require.Equal(t, sha256.Sum256(data), checksum)

ft := fullFT.Table()
require.Equal(t, dataSize/(2*megabyte), ft.NumFrames())
require.Equal(t, int64(dataSize), ft.UncompressedSize())

blob := anonymousGet(t, backend, key)
require.Equal(t, ft.CompressedSize(), int64(len(blob)))

decompressed, err := decompressAll(ft, blob)
require.NoError(t, err)
require.Equal(t, sha256.Sum256(data), sha256.Sum256(decompressed),
"stored blob must decompress back to the original")
})
}
}

// TestGCSXMLCompressedEmptyFile verifies the single empty part shipped for
// zero-byte inputs completes against a real XML multipart implementation.
func TestGCSXMLCompressedEmptyFile(t *testing.T) {
t.Parallel()

inputPath := writeTempFile(t, nil)
backend := gcsXMLBackend(t)
key := testKey("gcs-compressed-empty")

fullFT, checksum, err := storeFileCompressed(t.Context(), inputPath, testCompressConfig(), 4, PutOptions{},
func(metadata ObjectMetadata) (partUploader, error) {
return gcsXMLUploader(t, backend, key, metadata, nil), nil
})
require.NoError(t, err)
require.Equal(t, sha256.Sum256(nil), checksum)
require.Equal(t, 0, fullFT.Table().NumFrames())
require.Empty(t, anonymousGet(t, backend, key))
}

// TestGCSXMLRetryOnTransientFailure fails the first attempt of every request
// (initiate, each part, complete) with an injected 500 and verifies
// retryablehttp's retry with ReaderFunc body replay: the whole upload
// succeeds, retried parts re-send byte-identical bodies, and the stored blob
// decompresses cleanly.
func TestGCSXMLRetryOnTransientFailure(t *testing.T) {
t.Parallel()

const dataSize = 32 * megabyte
data := generateSemiRandomData(dataSize)
inputPath := writeTempFile(t, data)

backend := gcsXMLBackend(t)
key := testKey("gcs-retry-transient")
ft := newFaultInjectingTransport()

cfg := CompressConfig{
Enabled: true,
Type: CompressionLZ4.String(),
FrameSizeKB: 2 * 1024,
MinPartSizeMB: 5,
FrameEncodeWorkers: 4,
EncoderConcurrency: 1,
}

fullFT, checksum, err := storeFileCompressed(t.Context(), inputPath, cfg, 4, PutOptions{},
func(metadata ObjectMetadata) (partUploader, error) {
return gcsXMLUploader(t, backend, key, metadata, ft), nil
})
require.NoError(t, err, "upload must survive one injected 500 per request")
require.Equal(t, sha256.Sum256(data), checksum)

blob := anonymousGet(t, backend, key)
decompressed, err := decompressAll(fullFT.Table(), blob)
require.NoError(t, err)
require.Equal(t, sha256.Sum256(data), sha256.Sum256(decompressed),
"stored blob after injected faults must decompress to the original")

ft.mu.Lock()
defer ft.mu.Unlock()

require.NotEmpty(t, ft.partBodySizes)
require.Greater(t, ft.injected, len(ft.partBodySizes),
"should have injected faults beyond part uploads (initiate/complete)")
t.Logf("injected %d faults across %d distinct requests (%d parts)",
ft.injected, len(ft.seen), len(ft.partBodySizes))

for part, sizes := range ft.partBodySizes {
require.Len(t, sizes, 2, "part %s: expected exactly one failed and one successful attempt", part)
require.Positive(t, sizes[0], "part %s: injected attempt consumed no body", part)
require.Equal(t, sizes[0], sizes[1],
"part %s: retry sent %d bytes but first attempt sent %d — ReaderFunc body replay is broken",
part, sizes[1], sizes[0])
}
}

// TestGCSXMLUploadFileInParallel exercises the uncompressed >=50MB multipart
// path (UploadFileInParallel) against MinIO: fixed 50 MB chunks uploaded
// concurrently, with the overlapped checksum hasher.
func TestGCSXMLUploadFileInParallel(t *testing.T) {
t.Parallel()

const dataSize = 120 * megabyte // 3 parts at the fixed 50 MB chunk size
data := generateSemiRandomData(dataSize)
inputPath := writeTempFile(t, data)

backend := gcsXMLBackend(t)
key := testKey("gcs-parallel-upload")
up := gcsXMLUploader(t, backend, key, nil, nil)

hasher := sha256.New()
n, err := up.UploadFileInParallel(t.Context(), inputPath, 4, hasher)
require.NoError(t, err)
require.Equal(t, int64(dataSize), n)
require.Equal(t, sha256.Sum256(data), [32]byte(hasher.Sum(nil)),
"overlapped hasher must checksum the whole file")

got := anonymousGet(t, backend, key)
require.Equal(t, sha256.Sum256(data), sha256.Sum256(got))
}

// TestGCPUploadFileInParallelEmptyFile verifies the empty-file path ships its
// single zero-byte part with an explicit Content-Length rather than chunked
// transfer encoding, which S3-compatible XML backends (MinIO here) reject with
// 411 MissingContentLength.
func TestGCPUploadFileInParallelEmptyFile(t *testing.T) {
t.Parallel()

backend := gcsXMLBackend(t)
key := testKey("gcp-parallel-empty")
up := gcsXMLUploader(t, backend, key, nil, nil)

inputPath := writeTempFile(t, nil)
n, err := up.UploadFileInParallel(t.Context(), inputPath, 2, nil)
require.NoError(t, err, "empty file must upload without a 411 chunked-encoding rejection")
require.Zero(t, n)
require.Empty(t, anonymousGet(t, backend, key))
}
Loading
Loading