Skip to content

Commit e5aa742

Browse files
authored
Merge pull request #1472 from Altinity/hotfix_putMultipartCRC32
refactoring for CRC32 to properly size calculation,
2 parents 34d0669 + c8b3d90 commit e5aa742

3 files changed

Lines changed: 118 additions & 37 deletions

File tree

ChangeLog.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@ NEW FEATURES
55
- add `gcs.allow_multipart_download` (env `GCS_ALLOW_MULTIPART_DOWNLOAD`, default `false`) with `gcs.download_concurrency` — download each file as parallel range reads (part size is `chunk_size`) into a temporary file, mirrors `s3.allow_multipart_download`, fix [#1028](https://github.com/Altinity/clickhouse-backup/issues/1028)
66
- add `general.disable_environment_override` (settable ONLY in the config file, no environment variable name on purpose, default `false`) — when `true`, config values come only from the config file: all environment variables and the `--env` CLI flag are ignored during config loading; protects against accidental overrides such as Kubernetes service-discovery variables, fix [#1079](https://github.com/Altinity/clickhouse-backup/issues/1079)
77

8+
BUG FIXES
9+
- fix silent truncation of compressed backup uploads to S3 when `s3.check_sum_algorithm: CRC32` (Object Lock buckets): the multipart upload stopped at the declared size (sum of raw file sizes), dropping the tar headers/padding/trailer tail of the archive; the stream is now consumed until EOF, fix [#1471](https://github.com/Altinity/clickhouse-backup/issues/1471)
10+
811
# v2.7.4
912

1013
NEW FEATURES

pkg/storage/s3.go

Lines changed: 48 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -448,7 +448,7 @@ func (s *S3) PutFileAbsolute(ctx context.Context, key string, r io.ReadCloser, l
448448

449449
// transfermanager.UploadObject sends the part checksum as a trailer, which S3 Object Lock rejects on UploadPart. Fall back to manual multipart with per-part x-amz-checksum-crc32 header, fix https://github.com/Altinity/clickhouse-backup/issues/829
450450
if s.Config.CheckSumAlgorithm == string(s3types.ChecksumAlgorithmCrc32) && localSize > partSize {
451-
return s.putFileMultipartCRC32(ctx, &params, r, localSize, partSize)
451+
return s.putFileMultipartCRC32(ctx, s.client, &params, r, localSize, partSize)
452452
}
453453

454454
uploadInput := &transfermanager.UploadObjectInput{
@@ -480,15 +480,23 @@ func (s *S3) PutFileAbsolute(ctx context.Context, key string, r io.ReadCloser, l
480480
return nil
481481
}
482482

483-
func (s *S3) putFileMultipartCRC32(ctx context.Context, putParams *s3.PutObjectInput, r io.Reader, localSize, partSize int64) error {
483+
// s3MultipartAPI is the subset of *s3.Client used by putFileMultipartCRC32, extracted for unit testing
484+
type s3MultipartAPI interface {
485+
CreateMultipartUpload(ctx context.Context, params *s3.CreateMultipartUploadInput, optFns ...func(*s3.Options)) (*s3.CreateMultipartUploadOutput, error)
486+
UploadPart(ctx context.Context, params *s3.UploadPartInput, optFns ...func(*s3.Options)) (*s3.UploadPartOutput, error)
487+
CompleteMultipartUpload(ctx context.Context, params *s3.CompleteMultipartUploadInput, optFns ...func(*s3.Options)) (*s3.CompleteMultipartUploadOutput, error)
488+
AbortMultipartUpload(ctx context.Context, params *s3.AbortMultipartUploadInput, optFns ...func(*s3.Options)) (*s3.AbortMultipartUploadOutput, error)
489+
}
490+
491+
func (s *S3) putFileMultipartCRC32(ctx context.Context, s3MultipartAPIClient s3MultipartAPI, putParams *s3.PutObjectInput, r io.Reader, localSize, partSize int64) error {
484492
createParams := &s3.CreateMultipartUploadInput{
485493
Bucket: putParams.Bucket,
486494
Key: putParams.Key,
487495
StorageClass: putParams.StorageClass,
488496
}
489497
s.enrichCreateMultipartUploadParams(createParams)
490498

491-
initResp, err := s.client.CreateMultipartUpload(ctx, createParams)
499+
initResp, err := s3MultipartAPIClient.CreateMultipartUpload(ctx, createParams)
492500
if err != nil {
493501
return errors.Wrap(err, "S3 putFileMultipartCRC32 CreateMultipartUpload")
494502
}
@@ -503,54 +511,57 @@ func (s *S3) putFileMultipartCRC32(ctx context.Context, putParams *s3.PutObjectI
503511
if s.Config.RequestPayer != "" {
504512
abortParams.RequestPayer = s3types.RequestPayer(s.Config.RequestPayer)
505513
}
506-
if _, abortErr := s.client.AbortMultipartUpload(context.Background(), abortParams); abortErr != nil {
514+
if _, abortErr := s3MultipartAPIClient.AbortMultipartUpload(context.Background(), abortParams); abortErr != nil {
507515
return errors.Wrapf(cause, "aborting putFileMultipartCRC32 multipart upload: %v, original error was", abortErr)
508516
}
509517
return cause
510518
}
511519

520+
// localSize is only a capacity hint: for compressed streams UploadCompressedStream declares the
521+
// sum of raw file sizes, while the actual tar stream is larger (headers/padding/trailer), so read
522+
// until EOF instead of stopping at localSize, fix https://github.com/Altinity/clickhouse-backup/issues/1471
512523
buf := make([]byte, partSize)
513524
parts := make([]s3types.CompletedPart, 0, (localSize+partSize-1)/partSize)
514525
var partNumber int32 = 1
515-
remaining := localSize
516-
for remaining > 0 {
517-
toRead := partSize
518-
if remaining < toRead {
519-
toRead = remaining
520-
}
521-
if _, readErr := io.ReadFull(r, buf[:toRead]); readErr != nil {
526+
for {
527+
n, readErr := io.ReadFull(r, buf)
528+
if readErr != nil && readErr != io.EOF && readErr != io.ErrUnexpectedEOF {
522529
return abort(errors.Wrapf(readErr, "S3 putFileMultipartCRC32 read part=%d", partNumber))
523530
}
524-
h := crc32.NewIEEE()
525-
if _, writeErr := h.Write(buf[:toRead]); writeErr != nil {
526-
return errors.Wrapf(writeErr, "S3 putFileMultipartCRC32 write part=%d", partNumber)
527-
}
528-
uploadParams := &s3.UploadPartInput{
529-
Bucket: putParams.Bucket,
530-
Key: putParams.Key,
531-
UploadId: uploadID,
532-
PartNumber: aws.Int32(partNumber),
533-
Body: bytes.NewReader(buf[:toRead]),
534-
ChecksumAlgorithm: s3types.ChecksumAlgorithmCrc32,
535-
ChecksumCRC32: aws.String(base64.StdEncoding.EncodeToString(h.Sum(nil))),
536-
}
537-
if s.Config.RequestPayer != "" {
538-
uploadParams.RequestPayer = s3types.RequestPayer(s.Config.RequestPayer)
531+
if n > 0 {
532+
h := crc32.NewIEEE()
533+
if _, writeErr := h.Write(buf[:n]); writeErr != nil {
534+
return errors.Wrapf(writeErr, "S3 putFileMultipartCRC32 write part=%d", partNumber)
535+
}
536+
uploadParams := &s3.UploadPartInput{
537+
Bucket: putParams.Bucket,
538+
Key: putParams.Key,
539+
UploadId: uploadID,
540+
PartNumber: aws.Int32(partNumber),
541+
Body: bytes.NewReader(buf[:n]),
542+
ChecksumAlgorithm: s3types.ChecksumAlgorithmCrc32,
543+
ChecksumCRC32: aws.String(base64.StdEncoding.EncodeToString(h.Sum(nil))),
544+
}
545+
if s.Config.RequestPayer != "" {
546+
uploadParams.RequestPayer = s3types.RequestPayer(s.Config.RequestPayer)
547+
}
548+
partResp, uploadErr := s3MultipartAPIClient.UploadPart(ctx, uploadParams)
549+
if uploadErr != nil {
550+
return abort(errors.Wrapf(uploadErr, "S3 putFileMultipartCRC32 UploadPart part=%d", partNumber))
551+
}
552+
parts = append(parts, s3types.CompletedPart{
553+
ETag: partResp.ETag,
554+
PartNumber: aws.Int32(partNumber),
555+
ChecksumCRC32: partResp.ChecksumCRC32,
556+
})
557+
partNumber++
539558
}
540-
partResp, uploadErr := s.client.UploadPart(ctx, uploadParams)
541-
if uploadErr != nil {
542-
return abort(errors.Wrapf(uploadErr, "S3 putFileMultipartCRC32 UploadPart part=%d", partNumber))
559+
if readErr != nil {
560+
break
543561
}
544-
parts = append(parts, s3types.CompletedPart{
545-
ETag: partResp.ETag,
546-
PartNumber: aws.Int32(partNumber),
547-
ChecksumCRC32: partResp.ChecksumCRC32,
548-
})
549-
partNumber++
550-
remaining -= toRead
551562
}
552563

553-
if _, completeErr := s.client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{
564+
if _, completeErr := s3MultipartAPIClient.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{
554565
Bucket: putParams.Bucket,
555566
Key: putParams.Key,
556567
UploadId: uploadID,

pkg/storage/s3_test.go

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,80 @@
11
package storage
22

33
import (
4+
"bytes"
5+
"context"
46
"errors"
57
"fmt"
68
"testing"
79

10+
"github.com/Altinity/clickhouse-backup/v2/pkg/config"
11+
"github.com/aws/aws-sdk-go-v2/aws"
12+
"github.com/aws/aws-sdk-go-v2/service/s3"
813
"github.com/aws/smithy-go"
914
)
1015

16+
type fakeS3MultipartAPI struct {
17+
uploaded bytes.Buffer
18+
completed bool
19+
aborted bool
20+
}
21+
22+
func (f *fakeS3MultipartAPI) CreateMultipartUpload(_ context.Context, _ *s3.CreateMultipartUploadInput, _ ...func(*s3.Options)) (*s3.CreateMultipartUploadOutput, error) {
23+
return &s3.CreateMultipartUploadOutput{UploadId: aws.String("test-upload-id")}, nil
24+
}
25+
26+
func (f *fakeS3MultipartAPI) UploadPart(_ context.Context, params *s3.UploadPartInput, _ ...func(*s3.Options)) (*s3.UploadPartOutput, error) {
27+
if _, err := f.uploaded.ReadFrom(params.Body); err != nil {
28+
return nil, err
29+
}
30+
return &s3.UploadPartOutput{
31+
ETag: aws.String(fmt.Sprintf("etag-%d", *params.PartNumber)),
32+
ChecksumCRC32: params.ChecksumCRC32,
33+
}, nil
34+
}
35+
36+
func (f *fakeS3MultipartAPI) CompleteMultipartUpload(_ context.Context, _ *s3.CompleteMultipartUploadInput, _ ...func(*s3.Options)) (*s3.CompleteMultipartUploadOutput, error) {
37+
f.completed = true
38+
return &s3.CompleteMultipartUploadOutput{}, nil
39+
}
40+
41+
func (f *fakeS3MultipartAPI) AbortMultipartUpload(_ context.Context, _ *s3.AbortMultipartUploadInput, _ ...func(*s3.Options)) (*s3.AbortMultipartUploadOutput, error) {
42+
f.aborted = true
43+
return &s3.AbortMultipartUploadOutput{}, nil
44+
}
45+
46+
// TestPutFileMultipartCRC32ReadsUntilEOF - https://github.com/Altinity/clickhouse-backup/issues/1471
47+
// UploadCompressedStream declares localSize as the sum of raw on-disk file sizes, but the actual
48+
// tar stream is larger (512-byte header per file, 512-byte content padding, 1024-byte trailer).
49+
// The upload must consume the stream until EOF instead of stopping at the declared size,
50+
// otherwise the archive tail is silently truncated.
51+
func TestPutFileMultipartCRC32ReadsUntilEOF(t *testing.T) {
52+
const partSize = 8 * 1024
53+
const declaredSize = 3*partSize + 100 // sum of raw file sizes, what UploadCompressedStream computes
54+
const trueSize = declaredSize + 2560 // + tar headers/padding/trailer
55+
56+
data := make([]byte, trueSize)
57+
for i := range data {
58+
data[i] = byte(i % 251)
59+
}
60+
61+
fake := &fakeS3MultipartAPI{}
62+
s := &S3{Config: &config.S3Config{Bucket: "bucket"}}
63+
params := &s3.PutObjectInput{Bucket: aws.String("bucket"), Key: aws.String("backup/shadow/default/table/default_all_1_1_0.tar")}
64+
if err := s.putFileMultipartCRC32(context.Background(), fake, params, bytes.NewReader(data), declaredSize, partSize); err != nil {
65+
t.Fatalf("putFileMultipartCRC32 error: %v", err)
66+
}
67+
if !fake.completed {
68+
t.Fatal("CompleteMultipartUpload was not called")
69+
}
70+
if fake.aborted {
71+
t.Fatal("AbortMultipartUpload was called unexpectedly")
72+
}
73+
if !bytes.Equal(fake.uploaded.Bytes(), data) {
74+
t.Fatalf("uploaded %d bytes, want %d - the tail of the stream was truncated", fake.uploaded.Len(), trueSize)
75+
}
76+
}
77+
1178
// TestS3CopySource - `x-amz-copy-source` must be URL-encoded, keys with literal `%`/`#`/non-ASCII
1279
// (TablePathEncode names in shadow paths) were decoded server side into a different key, so
1380
// CopyObject failed with NoSuchKey while StatFile on the same key succeeded

0 commit comments

Comments
 (0)