From 6e3d2e0941da8a0b08514b6f446f694f39dfb433 Mon Sep 17 00:00:00 2001 From: sneak Date: Wed, 7 Oct 2026 09:18:31 +0000 Subject: [PATCH] Pass s3.part_size to the multipart uploader (closes #232) s3.part_size was loaded and defaulted but never reached the S3 client, whose uploader used a fixed 10MiB part. The client now takes the part size from the config, for storage_url and for the s3.* fields, and config load rejects a value below 5MiB or above 5GiB, the part sizes S3 accepts. The docs gave the default as 5MB, which the config file reads as 5,000,000 bytes, below the minimum; they now say 5MiB. Judgement call: the 5GiB maximum is enforced along with the 5MiB minimum the issue names. Trap: at the 5MiB default the uploader's 10,000-part limit caps one upload at about 48.8GiB, down from about 97.7GiB; blob_size_limit is not checked against it. Model: opus-5-5 --- TODO.md | 9 ++++ config.example.yml | 8 ++-- internal/cli/config.go | 2 +- internal/config/config.go | 12 ++++- internal/config/config_test.go | 71 +++++++++++++++++++++++++++ internal/s3/client.go | 12 +++-- internal/s3/module.go | 1 + internal/storage/module.go | 2 + internal/storage/s3_test.go | 87 ++++++++++++++++++++++++++++++++++ 9 files changed, 193 insertions(+), 11 deletions(-) diff --git a/TODO.md b/TODO.md index 0964f46..3ab0a3f 100644 --- a/TODO.md +++ b/TODO.md @@ -22,6 +22,15 @@ the tag exists and is exercised; what is left is merging `next` to # Completed Steps +- 2026-10-07: Made `s3.part_size` set the multipart upload part size + ([issue #232](https://git.eeqj.de/sneak/vaultik/issues/232)). It was + loaded and defaulted but never passed to the S3 client, whose uploader + used a fixed 10MiB part. It now reaches the uploader for `storage_url` + and for the `s3.*` fields, and a part size S3 refuses, below 5MiB or + above 5GiB, fails at config load. The docs gave the default as `5MB`, + which the config file reads as 5,000,000 bytes, below the minimum; they + now say `5MiB`. + - 2026-10-07: Made `remote info` stop reporting a snapshot's blobs as orphaned when its manifest cannot be read, and stop printing raw names from under `metadata/` diff --git a/config.example.yml b/config.example.yml index 6dc50de..db5c54e 100644 --- a/config.example.yml +++ b/config.example.yml @@ -287,10 +287,10 @@ storage_url: "rclone://myremote/path/to/backups" # #use_ssl: true # # # Part size for multipart uploads -# # Minimum 5MB, affects memory usage during upload -# # Supports: 5MB, 10M, 100MiB, etc. -# # Default: 5MB -# #part_size: 5MB +# # Minimum 5MiB, maximum 5GiB; affects memory usage during upload +# # Supports: 10MB, 16MiB, 100MiB, etc. (5MB is below the minimum) +# # Default: 5MiB +# #part_size: 5MiB # Path to local SQLite index database # This database tracks file state for incremental backups diff --git a/internal/cli/config.go b/internal/cli/config.go index b00a602..545f81a 100644 --- a/internal/cli/config.go +++ b/internal/cli/config.go @@ -205,7 +205,7 @@ storage_url: "" # access_key_id: YOUR_ACCESS_KEY # secret_access_key: YOUR_SECRET_KEY # # region: us-east-1 # Default: us-east-1 -# # part_size: 5MB # Multipart upload part size. Default: 5MB +# # part_size: 5MiB # Upload part size, 5MiB to 5GiB. Default: 5MiB # # For the s3:// form, disable TLS with ?ssl=false in the URL, not use_ssl. # ─── OPTIONAL ──────────────────────────────────────────────────────────────── diff --git a/internal/config/config.go b/internal/config/config.go index 1448a32..d3cc633 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -33,11 +33,14 @@ const secretKeyPrefix = "AGE-SECRET-KEY-" const ( defaultBlobSizeLimit = Size(10 * 1024 * 1024 * 1024) // 10GB defaultChunkSize = Size(10 * 1024 * 1024) // 10MB - defaultS3PartSize = Size(5 * 1024 * 1024) // 5MB + defaultS3PartSize = Size(5 * 1024 * 1024) // 5MiB defaultCompressionLevel = 3 minChunkSize = 1024 * 1024 // 1MB minCompressionLevel = 1 maxCompressionLevel = 19 + // S3 accepts a multipart upload part from 5MiB to 5GiB. + minS3PartSize = 5 * 1024 * 1024 + maxS3PartSize = 5 * 1024 * 1024 * 1024 ) // Sentinel validation errors. @@ -55,6 +58,7 @@ var ( "blob_size_limit must be at least the largest chunk the chunker can " + "emit (chunk_size times the FastCDC size spread)") errBadCompression = errors.New("compression_level must be between 1 and 19") + errBadS3PartSize = errors.New("s3.part_size must be between 5MiB and 5GiB") errBadStorageScheme = errors.New( "storage_url must start with s3://, file://, or rclone://") errStorageNotConfigured = errors.New( @@ -332,6 +336,7 @@ func Load(path string) (*Config, error) { // (chunk_size times chunker.ChunkSizeSpread), so a single-chunk blob never // exceeds the configured limit // - Compression level must be between 1 and 19 +// - S3 part size must be between 5MiB and 5GiB, the part sizes S3 accepts // // Returns an error describing the first validation failure encountered. func (c *Config) Validate() error { @@ -376,6 +381,11 @@ func (c *Config) Validate() error { return errBadCompression } + if c.S3.PartSize.Int64() < minS3PartSize || + c.S3.PartSize.Int64() > maxS3PartSize { + return errBadS3PartSize + } + return nil } diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 52b46ad..d4849e5 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -166,6 +166,7 @@ func TestValidateBlobSizeLimit(t *testing.T) { ChunkSize: chunkSize, BlobSizeLimit: blobLimit, CompressionLevel: 3, + S3: S3Config{PartSize: defaultS3PartSize}, } } @@ -221,6 +222,75 @@ func TestValidateBlobSizeLimit(t *testing.T) { } } +// TestValidateS3PartSize checks that s3.part_size is held to the part sizes +// S3 accepts, 5MiB to 5GiB. "5MB" in the config file is 5,000,000 bytes, +// below the minimum. +func TestValidateS3PartSize(t *testing.T) { + t.Parallel() + + newConfig := func(partSize Size) *Config { + return &Config{ + Snapshots: map[string]SnapshotConfig{"test": {Paths: []string{"/tmp/src"}}}, + StorageURL: "file:///tmp/vaultik-test-store", + ChunkSize: defaultChunkSize, + BlobSizeLimit: defaultBlobSizeLimit, + CompressionLevel: defaultCompressionLevel, + S3: S3Config{PartSize: partSize}, + } + } + + tests := []struct { + name string + partSize Size + wantErr bool + }{ + { + name: "5MB is rejected", + partSize: 5_000_000, + wantErr: true, + }, + { + name: "one byte below 5MiB is rejected", + partSize: minS3PartSize - 1, + wantErr: true, + }, + { + name: "5MiB is accepted", + partSize: minS3PartSize, + wantErr: false, + }, + { + name: "5GiB is accepted", + partSize: maxS3PartSize, + wantErr: false, + }, + { + name: "one byte above 5GiB is rejected", + partSize: maxS3PartSize + 1, + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + err := newConfig(tt.partSize).Validate() + if tt.wantErr { + if !errors.Is(err, errBadS3PartSize) { + t.Fatalf("Validate() error = %v, want errBadS3PartSize", err) + } + + return + } + + if err != nil { + t.Fatalf("Validate() unexpected error: %v", err) + } + }) + } +} + // TestValidateAgeRecipients checks that recipients are parsed at config load // (a bad entry fails immediately, not mid-backup) and that no invalid entry — // least of all a pasted secret key — is echoed in the error. An empty list @@ -236,6 +306,7 @@ func TestValidateAgeRecipients(t *testing.T) { ChunkSize: Size(10 * 1024 * 1024), BlobSizeLimit: Size(10 * 1024 * 1024 * 1024), CompressionLevel: 3, + S3: S3Config{PartSize: defaultS3PartSize}, } } diff --git a/internal/s3/client.go b/internal/s3/client.go index 43f8725..ae23054 100644 --- a/internal/s3/client.go +++ b/internal/s3/client.go @@ -27,10 +27,12 @@ type Client struct { bucket string prefix string endpoint string + partSize int64 } // Config contains S3 client configuration. -// All fields are required except Prefix, which defaults to an empty string. +// All fields are required except Prefix, which defaults to an empty string, +// and PartSize, where zero means the SDK default of 5 MiB. // A non-empty Prefix is joined to every key with one "/", whether or not // it ends with one. // The Endpoint field should include the protocol (http:// or https://). @@ -41,6 +43,8 @@ type Config struct { AccessKeyID string SecretAccessKey string Region string + // PartSize is the size in bytes of each part of a multipart upload. + PartSize int64 } // nopLogger is a logger that discards all output. @@ -90,6 +94,7 @@ func NewClient(ctx context.Context, cfg Config) (*Client, error) { bucket: cfg.Bucket, prefix: prefix, endpoint: cfg.Endpoint, + partSize: cfg.PartSize, }, nil } @@ -123,12 +128,9 @@ func (c *Client) PutObjectWithProgress( ) error { fullKey := c.prefix + key - // uploadPartSize is 10MB for better progress granularity. - const uploadPartSize = 10 * 1024 * 1024 - // Create an uploader with the S3 client uploader := manager.NewUploader(c.s3Client, func(u *manager.Uploader) { - u.PartSize = uploadPartSize + u.PartSize = c.partSize }) // Create a progress reader that tracks upload progress diff --git a/internal/s3/module.go b/internal/s3/module.go index 33b24ba..71f1b88 100644 --- a/internal/s3/module.go +++ b/internal/s3/module.go @@ -28,6 +28,7 @@ func provideClient(lc fx.Lifecycle, cfg *config.Config) (*Client, error) { AccessKeyID: cfg.S3.AccessKeyID, SecretAccessKey: cfg.S3.SecretAccessKey, Region: cfg.S3.Region, + PartSize: cfg.S3.PartSize.Int64(), }) if err != nil { return nil, err diff --git a/internal/storage/module.go b/internal/storage/module.go index 0e400fd..8d00ea1 100644 --- a/internal/storage/module.go +++ b/internal/storage/module.go @@ -99,6 +99,7 @@ func storerFromParsedS3URL(parsed *URL, cfg *config.Config) (Storer, error) { AccessKeyID: cfg.S3.AccessKeyID, SecretAccessKey: cfg.S3.SecretAccessKey, Region: region, + PartSize: cfg.S3.PartSize.Int64(), }) if err != nil { return nil, fmt.Errorf("creating S3 client: %w", err) @@ -134,6 +135,7 @@ func storerFromLegacyS3Config(cfg *config.Config) (Storer, error) { AccessKeyID: cfg.S3.AccessKeyID, SecretAccessKey: cfg.S3.SecretAccessKey, Region: region, + PartSize: cfg.S3.PartSize.Int64(), }) if err != nil { return nil, fmt.Errorf("creating S3 client: %w", err) diff --git a/internal/storage/s3_test.go b/internal/storage/s3_test.go index 2d581f7..f75be00 100644 --- a/internal/storage/s3_test.go +++ b/internal/storage/s3_test.go @@ -1,11 +1,14 @@ package storage_test import ( + "bytes" "context" "errors" + "net/http" "net/http/httptest" "slices" "strings" + "sync/atomic" "testing" "github.com/johannesboyne/gofakes3" @@ -171,6 +174,90 @@ func TestS3URLPrefixKeyLayout(t *testing.T) { } } +// TestS3UploadUsesConfiguredPartSize checks that s3.part_size reaches the +// multipart uploader, through storage_url and through the s3.* fields. An +// object three parts long must arrive as three parts; at the SDK's default +// of 5 MiB it would arrive as four. +func TestS3UploadUsesConfiguredPartSize(t *testing.T) { + t.Parallel() + + const ( + partSize = 6 * 1024 * 1024 + wantParts = 3 + ) + + backend := s3mem.New() + + err := backend.CreateBucket(s3TestBucket) + if err != nil { + t.Fatalf("create bucket: %v", err) + } + + // Every part of a multipart upload is one request with a partNumber. + var parts atomic.Int32 + + fake := gofakes3.New(backend).Server() + srv := httptest.NewServer(http.HandlerFunc( + func(w http.ResponseWriter, r *http.Request) { + if r.URL.Query().Has("partNumber") { + parts.Add(1) + } + + fake.ServeHTTP(w, r) + })) + t.Cleanup(srv.Close) + + cases := []struct { + name string + cfg *config.Config + }{ + { + name: "storage_url", + cfg: &config.Config{ + StorageURL: "s3://" + s3TestBucket + "?endpoint=" + srv.URL, + S3: config.S3Config{ + AccessKeyID: "key", + SecretAccessKey: "secret", + PartSize: partSize, + }, + }, + }, + { + name: "s3.endpoint", + cfg: &config.Config{ + S3: config.S3Config{ + Endpoint: srv.URL, + Bucket: s3TestBucket, + AccessKeyID: "key", + SecretAccessKey: "secret", + PartSize: partSize, + }, + }, + }, + } + + for _, tc := range cases { + parts.Store(0) + + storer, err := storage.NewStorer(tc.cfg) + if err != nil { + t.Fatalf("%s: NewStorer: %v", tc.name, err) + } + + data := bytes.NewReader(make([]byte, wantParts*partSize)) + + err = storer.PutWithProgress( + context.Background(), "blob", data, data.Size(), nil) + if err != nil { + t.Fatalf("%s: PutWithProgress: %v", tc.name, err) + } + + if got := parts.Load(); got != wantParts { + t.Errorf("%s: uploaded in %d parts, want %d", tc.name, got, wantParts) + } + } +} + // listStreamKeys returns the keys ListStream yields under a prefix, and // fails the test on a listing error. func listStreamKeys(t *testing.T, s storage.Storer, prefix string) []string {