Pass s3.part_size to the multipart uploader (closes #232)
check / check (push) Waiting to run
check / check (push) Waiting to run
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, an explicit 0 included. A blob too large for S3's limit of 10,000 parts at that size is uploaded in larger parts, since the uploader cannot learn the size of the reader it is given. 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. Model: opus-5-5
This commit was merged in pull request #262.
This commit is contained in:
@@ -22,6 +22,16 @@ the tag exists and is exercised; what is left is merging `next` to
|
|||||||
|
|
||||||
# Completed Steps
|
# 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, `0` included, fails at config load. A blob too large for
|
||||||
|
S3's limit of 10,000 parts at the configured size is uploaded in larger
|
||||||
|
parts. 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 per-name retention work when the hostname contains `_`
|
- 2026-10-07: Made per-name retention work when the hostname contains `_`
|
||||||
([issue #230](https://git.eeqj.de/sneak/vaultik/issues/230)). A
|
([issue #230](https://git.eeqj.de/sneak/vaultik/issues/230)). A
|
||||||
snapshot ID is `hostname_name_timestamp`, and the name was read as
|
snapshot ID is `hostname_name_timestamp`, and the name was read as
|
||||||
|
|||||||
+5
-4
@@ -287,10 +287,11 @@ storage_url: "rclone://myremote/path/to/backups"
|
|||||||
# #use_ssl: true
|
# #use_ssl: true
|
||||||
#
|
#
|
||||||
# # Part size for multipart uploads
|
# # Part size for multipart uploads
|
||||||
# # Minimum 5MB, affects memory usage during upload
|
# # Minimum 5MiB, maximum 5GiB; affects memory usage during upload
|
||||||
# # Supports: 5MB, 10M, 100MiB, etc.
|
# # A blob too large for 10,000 parts of this size gets larger parts
|
||||||
# # Default: 5MB
|
# # Supports: 10MB, 16MiB, 100MiB, etc. (5MB is below the minimum)
|
||||||
# #part_size: 5MB
|
# # Default: 5MiB
|
||||||
|
# #part_size: 5MiB
|
||||||
|
|
||||||
# Path to local SQLite index database
|
# Path to local SQLite index database
|
||||||
# This database tracks file state for incremental backups
|
# This database tracks file state for incremental backups
|
||||||
|
|||||||
@@ -205,7 +205,7 @@ storage_url: ""
|
|||||||
# access_key_id: YOUR_ACCESS_KEY
|
# access_key_id: YOUR_ACCESS_KEY
|
||||||
# secret_access_key: YOUR_SECRET_KEY
|
# secret_access_key: YOUR_SECRET_KEY
|
||||||
# # region: us-east-1 # Default: us-east-1
|
# # 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.
|
# # For the s3:// form, disable TLS with ?ssl=false in the URL, not use_ssl.
|
||||||
|
|
||||||
# ─── OPTIONAL ────────────────────────────────────────────────────────────────
|
# ─── OPTIONAL ────────────────────────────────────────────────────────────────
|
||||||
|
|||||||
@@ -33,11 +33,14 @@ const secretKeyPrefix = "AGE-SECRET-KEY-"
|
|||||||
const (
|
const (
|
||||||
defaultBlobSizeLimit = Size(10 * 1024 * 1024 * 1024) // 10GB
|
defaultBlobSizeLimit = Size(10 * 1024 * 1024 * 1024) // 10GB
|
||||||
defaultChunkSize = Size(10 * 1024 * 1024) // 10MB
|
defaultChunkSize = Size(10 * 1024 * 1024) // 10MB
|
||||||
defaultS3PartSize = Size(5 * 1024 * 1024) // 5MB
|
defaultS3PartSize = Size(5 * 1024 * 1024) // 5MiB
|
||||||
defaultCompressionLevel = 3
|
defaultCompressionLevel = 3
|
||||||
minChunkSize = 1024 * 1024 // 1MB
|
minChunkSize = 1024 * 1024 // 1MB
|
||||||
minCompressionLevel = 1
|
minCompressionLevel = 1
|
||||||
maxCompressionLevel = 19
|
maxCompressionLevel = 19
|
||||||
|
// S3 accepts a multipart upload part from 5MiB to 5GiB.
|
||||||
|
minS3PartSize = 5 * 1024 * 1024
|
||||||
|
maxS3PartSize = 5 * 1024 * 1024 * 1024
|
||||||
)
|
)
|
||||||
|
|
||||||
// Sentinel validation errors.
|
// Sentinel validation errors.
|
||||||
@@ -55,6 +58,7 @@ var (
|
|||||||
"blob_size_limit must be at least the largest chunk the chunker can " +
|
"blob_size_limit must be at least the largest chunk the chunker can " +
|
||||||
"emit (chunk_size times the FastCDC size spread)")
|
"emit (chunk_size times the FastCDC size spread)")
|
||||||
errBadCompression = errors.New("compression_level must be between 1 and 19")
|
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(
|
errBadStorageScheme = errors.New(
|
||||||
"storage_url must start with s3://, file://, or rclone://")
|
"storage_url must start with s3://, file://, or rclone://")
|
||||||
errStorageNotConfigured = errors.New(
|
errStorageNotConfigured = errors.New(
|
||||||
@@ -244,6 +248,7 @@ func Load(path string) (*Config, error) {
|
|||||||
ChunkSize: defaultChunkSize,
|
ChunkSize: defaultChunkSize,
|
||||||
IndexPath: filepath.Join(xdg.DataHome, appName, "index.sqlite"),
|
IndexPath: filepath.Join(xdg.DataHome, appName, "index.sqlite"),
|
||||||
CompressionLevel: defaultCompressionLevel,
|
CompressionLevel: defaultCompressionLevel,
|
||||||
|
S3: S3Config{PartSize: defaultS3PartSize},
|
||||||
}
|
}
|
||||||
|
|
||||||
// Convert smartconfig data to YAML then unmarshal
|
// Convert smartconfig data to YAML then unmarshal
|
||||||
@@ -294,10 +299,6 @@ func Load(path string) (*Config, error) {
|
|||||||
cfg.S3.Region = "us-east-1"
|
cfg.S3.Region = "us-east-1"
|
||||||
}
|
}
|
||||||
|
|
||||||
if cfg.S3.PartSize == 0 {
|
|
||||||
cfg.S3.PartSize = defaultS3PartSize
|
|
||||||
}
|
|
||||||
|
|
||||||
// Check config file permissions (warn if world or group readable)
|
// Check config file permissions (warn if world or group readable)
|
||||||
//nolint:gosec // G703: config path is operator-supplied by design
|
//nolint:gosec // G703: config path is operator-supplied by design
|
||||||
info, statErr := os.Stat(path)
|
info, statErr := os.Stat(path)
|
||||||
@@ -332,6 +333,7 @@ func Load(path string) (*Config, error) {
|
|||||||
// (chunk_size times chunker.ChunkSizeSpread), so a single-chunk blob never
|
// (chunk_size times chunker.ChunkSizeSpread), so a single-chunk blob never
|
||||||
// exceeds the configured limit
|
// exceeds the configured limit
|
||||||
// - Compression level must be between 1 and 19
|
// - 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.
|
// Returns an error describing the first validation failure encountered.
|
||||||
func (c *Config) Validate() error {
|
func (c *Config) Validate() error {
|
||||||
@@ -376,6 +378,11 @@ func (c *Config) Validate() error {
|
|||||||
return errBadCompression
|
return errBadCompression
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if c.S3.PartSize.Int64() < minS3PartSize ||
|
||||||
|
c.S3.PartSize.Int64() > maxS3PartSize {
|
||||||
|
return errBadS3PartSize
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -166,6 +166,7 @@ func TestValidateBlobSizeLimit(t *testing.T) {
|
|||||||
ChunkSize: chunkSize,
|
ChunkSize: chunkSize,
|
||||||
BlobSizeLimit: blobLimit,
|
BlobSizeLimit: blobLimit,
|
||||||
CompressionLevel: 3,
|
CompressionLevel: 3,
|
||||||
|
S3: S3Config{PartSize: defaultS3PartSize},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -221,6 +222,120 @@ func TestValidateBlobSizeLimit(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestValidateS3PartSize checks that s3.part_size is held to the part sizes
|
||||||
|
// S3 accepts, 5MiB to 5GiB, by changing only the part size of the test
|
||||||
|
// config. "5MB" in the config file is 5,000,000 bytes, below the minimum.
|
||||||
|
func TestValidateS3PartSize(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
base, err := Load(os.Getenv("VAULTIK_CONFIG"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Failed to load config: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
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()
|
||||||
|
|
||||||
|
cfg := *base
|
||||||
|
cfg.S3.PartSize = tt.partSize
|
||||||
|
|
||||||
|
err := cfg.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)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestLoadS3PartSize checks that a config file without s3.part_size loads
|
||||||
|
// with the 5MiB default, and that an explicit 0 fails at load like any other
|
||||||
|
// part size S3 refuses.
|
||||||
|
func TestLoadS3PartSize(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
const withoutPartSize = "snapshots:\n" +
|
||||||
|
" test:\n" +
|
||||||
|
" paths: [/tmp/vaultik-test-source]\n" +
|
||||||
|
"storage_url: file:///tmp/vaultik-test-storage\n"
|
||||||
|
|
||||||
|
writeConfig := func(t *testing.T, text string) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
path := filepath.Join(t.TempDir(), "config.yml")
|
||||||
|
|
||||||
|
err := os.WriteFile(path, []byte(text), 0o600)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("write config: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return path
|
||||||
|
}
|
||||||
|
|
||||||
|
t.Run("absent loads as 5MiB", func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
cfg, err := Load(writeConfig(t, withoutPartSize))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Load() unexpected error: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if cfg.S3.PartSize != defaultS3PartSize {
|
||||||
|
t.Errorf("s3.part_size = %d, want %d",
|
||||||
|
cfg.S3.PartSize, defaultS3PartSize)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("0 is rejected", func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
_, err := Load(writeConfig(t, withoutPartSize+"s3:\n part_size: 0\n"))
|
||||||
|
if !errors.Is(err, errBadS3PartSize) {
|
||||||
|
t.Fatalf("Load() error = %v, want errBadS3PartSize", err)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// TestValidateAgeRecipients checks that recipients are parsed at config load
|
// TestValidateAgeRecipients checks that recipients are parsed at config load
|
||||||
// (a bad entry fails immediately, not mid-backup) and that no invalid entry —
|
// (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
|
// least of all a pasted secret key — is echoed in the error. An empty list
|
||||||
@@ -236,6 +351,7 @@ func TestValidateAgeRecipients(t *testing.T) {
|
|||||||
ChunkSize: Size(10 * 1024 * 1024),
|
ChunkSize: Size(10 * 1024 * 1024),
|
||||||
BlobSizeLimit: Size(10 * 1024 * 1024 * 1024),
|
BlobSizeLimit: Size(10 * 1024 * 1024 * 1024),
|
||||||
CompressionLevel: 3,
|
CompressionLevel: 3,
|
||||||
|
S3: S3Config{PartSize: defaultS3PartSize},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+23
-5
@@ -27,10 +27,12 @@ type Client struct {
|
|||||||
bucket string
|
bucket string
|
||||||
prefix string
|
prefix string
|
||||||
endpoint string
|
endpoint string
|
||||||
|
partSize int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// Config contains S3 client configuration.
|
// 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
|
// A non-empty Prefix is joined to every key with one "/", whether or not
|
||||||
// it ends with one.
|
// it ends with one.
|
||||||
// The Endpoint field should include the protocol (http:// or https://).
|
// The Endpoint field should include the protocol (http:// or https://).
|
||||||
@@ -41,6 +43,9 @@ type Config struct {
|
|||||||
AccessKeyID string
|
AccessKeyID string
|
||||||
SecretAccessKey string
|
SecretAccessKey string
|
||||||
Region string
|
Region string
|
||||||
|
// PartSize is the size in bytes of each part of a multipart upload.
|
||||||
|
// An upload too large for S3's limit of 10,000 parts gets larger parts.
|
||||||
|
PartSize int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// nopLogger is a logger that discards all output.
|
// nopLogger is a logger that discards all output.
|
||||||
@@ -90,6 +95,7 @@ func NewClient(ctx context.Context, cfg Config) (*Client, error) {
|
|||||||
bucket: cfg.Bucket,
|
bucket: cfg.Bucket,
|
||||||
prefix: prefix,
|
prefix: prefix,
|
||||||
endpoint: cfg.Endpoint,
|
endpoint: cfg.Endpoint,
|
||||||
|
partSize: cfg.PartSize,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -123,12 +129,9 @@ func (c *Client) PutObjectWithProgress(
|
|||||||
) error {
|
) error {
|
||||||
fullKey := c.prefix + key
|
fullKey := c.prefix + key
|
||||||
|
|
||||||
// uploadPartSize is 10MB for better progress granularity.
|
|
||||||
const uploadPartSize = 10 * 1024 * 1024
|
|
||||||
|
|
||||||
// Create an uploader with the S3 client
|
// Create an uploader with the S3 client
|
||||||
uploader := manager.NewUploader(c.s3Client, func(u *manager.Uploader) {
|
uploader := manager.NewUploader(c.s3Client, func(u *manager.Uploader) {
|
||||||
u.PartSize = uploadPartSize
|
u.PartSize = uploadPartSize(c.partSize, size)
|
||||||
})
|
})
|
||||||
|
|
||||||
// Create a progress reader that tracks upload progress
|
// Create a progress reader that tracks upload progress
|
||||||
@@ -149,6 +152,21 @@ func (c *Client) PutObjectWithProgress(
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// uploadPartSize returns the part size for an upload of size bytes: the
|
||||||
|
// configured part size (the SDK default when zero), raised where needed so
|
||||||
|
// the upload fits in S3's limit of 10,000 parts. The uploader cannot raise
|
||||||
|
// it itself, because it cannot seek the progress reader to learn its size.
|
||||||
|
func uploadPartSize(configured, size int64) int64 {
|
||||||
|
if configured == 0 {
|
||||||
|
configured = manager.DefaultUploadPartSize
|
||||||
|
}
|
||||||
|
|
||||||
|
maxParts := int64(manager.MaxUploadParts)
|
||||||
|
smallestThatFits := (size + maxParts - 1) / maxParts // rounded up
|
||||||
|
|
||||||
|
return max(configured, smallestThatFits)
|
||||||
|
}
|
||||||
|
|
||||||
// GetObject downloads an object from S3 with the specified key.
|
// GetObject downloads an object from S3 with the specified key.
|
||||||
// The key is automatically prefixed with the configured prefix.
|
// The key is automatically prefixed with the configured prefix.
|
||||||
// Returns a ReadCloser containing the object data. The caller must
|
// Returns a ReadCloser containing the object data. The caller must
|
||||||
|
|||||||
@@ -0,0 +1,55 @@
|
|||||||
|
package s3
|
||||||
|
|
||||||
|
import "testing"
|
||||||
|
|
||||||
|
// TestUploadPartSize checks that an upload too large for 10,000 parts of the
|
||||||
|
// configured size gets parts just large enough to fit in 10,000.
|
||||||
|
func TestUploadPartSize(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
const mib = 1024 * 1024
|
||||||
|
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
configured int64
|
||||||
|
size int64
|
||||||
|
want int64
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "an upload that fits keeps the configured size",
|
||||||
|
configured: 5 * mib,
|
||||||
|
size: 10 * 1024 * mib,
|
||||||
|
want: 5 * mib,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "exactly 10,000 parts keeps the configured size",
|
||||||
|
configured: 6 * mib,
|
||||||
|
size: 10_000 * 6 * mib,
|
||||||
|
want: 6 * mib,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "one byte more than 10,000 parts adds a byte to each",
|
||||||
|
configured: 6 * mib,
|
||||||
|
size: 10_000*6*mib + 1,
|
||||||
|
want: 6*mib + 1,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "zero means the SDK default of 5MiB",
|
||||||
|
configured: 0,
|
||||||
|
size: 1,
|
||||||
|
want: 5 * mib,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
got := uploadPartSize(tt.configured, tt.size)
|
||||||
|
if got != tt.want {
|
||||||
|
t.Errorf("uploadPartSize(%d, %d) = %d, want %d",
|
||||||
|
tt.configured, tt.size, got, tt.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -28,6 +28,7 @@ func provideClient(lc fx.Lifecycle, cfg *config.Config) (*Client, error) {
|
|||||||
AccessKeyID: cfg.S3.AccessKeyID,
|
AccessKeyID: cfg.S3.AccessKeyID,
|
||||||
SecretAccessKey: cfg.S3.SecretAccessKey,
|
SecretAccessKey: cfg.S3.SecretAccessKey,
|
||||||
Region: cfg.S3.Region,
|
Region: cfg.S3.Region,
|
||||||
|
PartSize: cfg.S3.PartSize.Int64(),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
@@ -99,6 +99,7 @@ func storerFromParsedS3URL(parsed *URL, cfg *config.Config) (Storer, error) {
|
|||||||
AccessKeyID: cfg.S3.AccessKeyID,
|
AccessKeyID: cfg.S3.AccessKeyID,
|
||||||
SecretAccessKey: cfg.S3.SecretAccessKey,
|
SecretAccessKey: cfg.S3.SecretAccessKey,
|
||||||
Region: region,
|
Region: region,
|
||||||
|
PartSize: cfg.S3.PartSize.Int64(),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("creating S3 client: %w", err)
|
return nil, fmt.Errorf("creating S3 client: %w", err)
|
||||||
@@ -134,6 +135,7 @@ func storerFromLegacyS3Config(cfg *config.Config) (Storer, error) {
|
|||||||
AccessKeyID: cfg.S3.AccessKeyID,
|
AccessKeyID: cfg.S3.AccessKeyID,
|
||||||
SecretAccessKey: cfg.S3.SecretAccessKey,
|
SecretAccessKey: cfg.S3.SecretAccessKey,
|
||||||
Region: region,
|
Region: region,
|
||||||
|
PartSize: cfg.S3.PartSize.Int64(),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("creating S3 client: %w", err)
|
return nil, fmt.Errorf("creating S3 client: %w", err)
|
||||||
|
|||||||
@@ -1,11 +1,14 @@
|
|||||||
package storage_test
|
package storage_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"slices"
|
"slices"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/johannesboyne/gofakes3"
|
"github.com/johannesboyne/gofakes3"
|
||||||
@@ -19,6 +22,13 @@ import (
|
|||||||
// s3TestBucket is the bucket created for each in-process S3 server.
|
// s3TestBucket is the bucket created for each in-process S3 server.
|
||||||
const s3TestBucket = "test-bucket"
|
const s3TestBucket = "test-bucket"
|
||||||
|
|
||||||
|
// Credentials for the tests that build a storer from a config.Config. The
|
||||||
|
// in-process S3 server accepts any.
|
||||||
|
const (
|
||||||
|
s3TestAccessKeyID = "key"
|
||||||
|
s3TestSecretAccessKey = "secret"
|
||||||
|
)
|
||||||
|
|
||||||
// newS3Storer builds an s3:// backend backed by a fresh in-process
|
// newS3Storer builds an s3:// backend backed by a fresh in-process
|
||||||
// S3 server. It reuses the same in-memory S3 harness (gofakes3 + s3mem
|
// S3 server. It reuses the same in-memory S3 harness (gofakes3 + s3mem
|
||||||
// over httptest) that internal/s3 and the not-found regression test use,
|
// over httptest) that internal/s3 and the not-found regression test use,
|
||||||
@@ -128,8 +138,8 @@ func TestS3URLPrefixKeyLayout(t *testing.T) {
|
|||||||
storer, err := storage.NewStorer(&config.Config{
|
storer, err := storage.NewStorer(&config.Config{
|
||||||
StorageURL: storageURL + "?endpoint=" + srv.URL,
|
StorageURL: storageURL + "?endpoint=" + srv.URL,
|
||||||
S3: config.S3Config{
|
S3: config.S3Config{
|
||||||
AccessKeyID: "key",
|
AccessKeyID: s3TestAccessKeyID,
|
||||||
SecretAccessKey: "secret",
|
SecretAccessKey: s3TestSecretAccessKey,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -171,6 +181,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: s3TestAccessKeyID,
|
||||||
|
SecretAccessKey: s3TestSecretAccessKey,
|
||||||
|
PartSize: partSize,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "s3.endpoint",
|
||||||
|
cfg: &config.Config{
|
||||||
|
S3: config.S3Config{
|
||||||
|
Endpoint: srv.URL,
|
||||||
|
Bucket: s3TestBucket,
|
||||||
|
AccessKeyID: s3TestAccessKeyID,
|
||||||
|
SecretAccessKey: s3TestSecretAccessKey,
|
||||||
|
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
|
// listStreamKeys returns the keys ListStream yields under a prefix, and
|
||||||
// fails the test on a listing error.
|
// fails the test on a listing error.
|
||||||
func listStreamKeys(t *testing.T, s storage.Storer, prefix string) []string {
|
func listStreamKeys(t *testing.T, s storage.Storer, prefix string) []string {
|
||||||
|
|||||||
+1
-1
@@ -19,7 +19,7 @@ s3:
|
|||||||
secret_access_key: test-secret-key
|
secret_access_key: test-secret-key
|
||||||
region: us-east-1
|
region: us-east-1
|
||||||
use_ssl: true
|
use_ssl: true
|
||||||
part_size: 5242880 # 5MB
|
part_size: 5242880 # 5MiB
|
||||||
index_path: /tmp/vaultik-test.sqlite
|
index_path: /tmp/vaultik-test.sqlite
|
||||||
chunk_size: 10MB
|
chunk_size: 10MB
|
||||||
blob_size_limit: 10GB
|
blob_size_limit: 10GB
|
||||||
|
|||||||
Reference in New Issue
Block a user