check / check (push) Waiting to run
The file and rclone listings skip an object whose name ends in `.partial`, the temporary name a `file://` or rclone upload writes before moving the object into place. `remote nuke` deletes only what the listings return, so it left the `.partial` objects killed uploads leave behind and still reported the destination store empty. Storer gains DeletePartialUploads. The file and rclone backends remove every `.partial` object under the prefix; S3 has none to remove, since it shows an object only once its upload completes. `remote nuke` calls it for `metadata/` and `blobs/` as its last step. Empty directories under a `file://` destination are still left behind. Model: opus-5-5
115 lines
2.8 KiB
Go
115 lines
2.8 KiB
Go
package storage
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
|
|
"sneak.berlin/go/vaultik/internal/s3"
|
|
)
|
|
|
|
// S3Storer wraps the existing s3.Client to implement Storer.
|
|
type S3Storer struct {
|
|
client *s3.Client
|
|
}
|
|
|
|
// NewS3Storer creates a new S3 storage backend.
|
|
func NewS3Storer(client *s3.Client) *S3Storer {
|
|
return &S3Storer{client: client}
|
|
}
|
|
|
|
// Put stores data at the specified key.
|
|
func (s *S3Storer) Put(ctx context.Context, key string, data io.Reader) error {
|
|
return s.client.PutObject(ctx, key, data)
|
|
}
|
|
|
|
// PutWithProgress stores data with progress reporting.
|
|
func (s *S3Storer) PutWithProgress(
|
|
ctx context.Context, key string, data io.Reader,
|
|
size int64, progress ProgressCallback,
|
|
) error {
|
|
// Convert storage.ProgressCallback to s3.ProgressCallback
|
|
var s3Progress s3.ProgressCallback
|
|
if progress != nil {
|
|
s3Progress = s3.ProgressCallback(progress)
|
|
}
|
|
|
|
return s.client.PutObjectWithProgress(ctx, key, data, size, s3Progress)
|
|
}
|
|
|
|
// Get retrieves data from the specified key.
|
|
// Returns ErrNotFound if the object does not exist.
|
|
func (s *S3Storer) Get(ctx context.Context, key string) (io.ReadCloser, error) {
|
|
rc, err := s.client.GetObject(ctx, key)
|
|
if err != nil {
|
|
if s3.IsNotFound(err) {
|
|
return nil, fmt.Errorf("get %q: %w", key, ErrNotFound)
|
|
}
|
|
|
|
return nil, err
|
|
}
|
|
|
|
return rc, nil
|
|
}
|
|
|
|
// Stat returns metadata about an object without retrieving its contents.
|
|
// Returns ErrNotFound if the object does not exist.
|
|
func (s *S3Storer) Stat(ctx context.Context, key string) (*ObjectInfo, error) {
|
|
info, err := s.client.StatObject(ctx, key)
|
|
if err != nil {
|
|
if s3.IsNotFound(err) {
|
|
return nil, fmt.Errorf("stat %q: %w", key, ErrNotFound)
|
|
}
|
|
|
|
return nil, err
|
|
}
|
|
|
|
return &ObjectInfo{
|
|
Key: info.Key,
|
|
Size: info.Size,
|
|
}, nil
|
|
}
|
|
|
|
// Delete removes an object.
|
|
func (s *S3Storer) Delete(ctx context.Context, key string) error {
|
|
return s.client.DeleteObject(ctx, key)
|
|
}
|
|
|
|
// List returns all keys with the given prefix.
|
|
func (s *S3Storer) List(ctx context.Context, prefix string) ([]string, error) {
|
|
return s.client.ListObjects(ctx, prefix)
|
|
}
|
|
|
|
// ListStream returns a channel of ObjectInfo for large result sets.
|
|
func (s *S3Storer) ListStream(ctx context.Context, prefix string) <-chan ObjectInfo {
|
|
ch := make(chan ObjectInfo)
|
|
|
|
go func() {
|
|
defer close(ch)
|
|
|
|
for info := range s.client.ListObjectsStream(ctx, prefix, false) {
|
|
ch <- ObjectInfo{
|
|
Key: info.Key,
|
|
Size: info.Size,
|
|
Err: info.Err,
|
|
}
|
|
}
|
|
}()
|
|
|
|
return ch
|
|
}
|
|
|
|
// DeletePartialUploads has nothing to remove: S3 shows an object only once
|
|
// its upload has completed, so an upload cut off part-way leaves none.
|
|
func (s *S3Storer) DeletePartialUploads(_ context.Context, _ string) error {
|
|
return nil
|
|
}
|
|
|
|
// Info returns human-readable storage location information.
|
|
func (s *S3Storer) Info() Info {
|
|
return Info{
|
|
Type: "s3",
|
|
Location: fmt.Sprintf("%s/%s", s.client.Endpoint(), s.client.BucketName()),
|
|
}
|
|
}
|