feat: add utils command to convert preexisting posix datasets

Closes #2304

Adds `versitygw utils convert-posix-dataset` (alias `cpd`), which makes a posix dataset not created by the gateway fully compatible with the posix backend. Every top level directory is treated as a bucket and gets a private ACL owned by `--access-key-id` (defaults to the root `--access`) and `BucketOwnerEnforced` object ownership; root level files are ignored. Metadata is stored in xattrs, or in `--sidecar-dir` when set.

`--calculate-etag` and `--checksum-algorithm` optionally compute the MD5 ETag and the `FULL_OBJECT` checksum of every object in a single streamed read. Objects are converted concurrently, tunable with `--concurrency` and `--read-buffer-size`. Existing bucket and object metadata is never overwritten.
This commit is contained in:
niksis02
2026-09-14 20:53:08 +04:00
parent f002073a07
commit 13b6bb061c
4 changed files with 591 additions and 61 deletions
+4 -33
View File
@@ -2007,35 +2007,6 @@ func getPartChecksum(algo types.ChecksumAlgorithm, part types.CompletedPart) str
}
}
func setStoredChecksum(checksum *s3response.Checksum, algo types.ChecksumAlgorithm, sum *string) {
if sum == nil {
return
}
switch algo {
case types.ChecksumAlgorithmCrc32:
checksum.CRC32 = sum
case types.ChecksumAlgorithmCrc32c:
checksum.CRC32C = sum
case types.ChecksumAlgorithmSha1:
checksum.SHA1 = sum
case types.ChecksumAlgorithmSha256:
checksum.SHA256 = sum
case types.ChecksumAlgorithmCrc64nvme:
checksum.CRC64NVME = sum
case types.ChecksumAlgorithmSha512:
checksum.SHA512 = sum
case types.ChecksumAlgorithmMd5:
checksum.MD5 = sum
case types.ChecksumAlgorithmXxhash64:
checksum.XXHASH64 = sum
case types.ChecksumAlgorithmXxhash3:
checksum.XXHASH3 = sum
case types.ChecksumAlgorithmXxhash128:
checksum.XXHASH128 = sum
}
}
func setUploadPartChecksum(res *s3.UploadPartOutput, algo types.ChecksumAlgorithm, sum *string) {
if sum == nil {
return
@@ -3715,7 +3686,7 @@ func (p *Posix) UploadPartWithPostFunc(ctx context.Context, input *s3.UploadPart
sum = hashRdr.Sum()
}
setStoredChecksum(&checksum, checksums.Algorithm, &sum)
checksum.SetSum(checksums.Algorithm, &sum)
setUploadPartChecksum(res, checksums.Algorithm, &sum)
err := p.storeChecksums(f.File(), bucket, partPath, checksum)
@@ -4028,7 +3999,7 @@ func (p *Posix) UploadPartCopy(ctx context.Context, upi *s3.UploadPartCopyInput)
}
sum := hashRdr.Sum()
setStoredChecksum(&checksums, algo, &sum)
checksums.SetSum(algo, &sum)
err := p.storeChecksums(f.File(), *upi.Bucket, partPath, checksums)
if err != nil {
@@ -4345,7 +4316,7 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje
Algorithm: checksumAlgorithm,
}
setStoredChecksum(&checksum, checksumAlgorithm, &expectedSum)
checksum.SetSum(checksumAlgorithm, &expectedSum)
err = p.storeChecksums(nil, *po.Bucket, *po.Key, checksum)
if err != nil {
@@ -4557,7 +4528,7 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje
Algorithm: checksumAlgorithm,
}
setStoredChecksum(&checksum, checksumAlgorithm, &sum)
checksum.SetSum(checksumAlgorithm, &sum)
err = p.storeChecksums(f.File(), *po.Bucket, *po.Key, checksum)
if err != nil {
return s3response.PutObjectOutput{}, fmt.Errorf("store checksum: %w", err)
+523
View File
@@ -15,21 +15,41 @@
package gwcli
import (
"context"
"crypto/md5"
"encoding/json"
"errors"
"fmt"
"hash"
"io"
"io/fs"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"sync/atomic"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/urfave/cli/v2"
"github.com/versity/versitygw/auth"
"github.com/versity/versitygw/backend"
"github.com/versity/versitygw/backend/meta"
"github.com/versity/versitygw/s3api/utils"
"github.com/versity/versitygw/s3err"
"github.com/versity/versitygw/s3event"
"github.com/versity/versitygw/s3response"
)
// UtilsCommand returns the "utils" subcommand, common to all versitygw
// binaries.
func UtilsCommand() *cli.Command {
algoValues := types.ChecksumAlgorithmCrc32.Values()
algos := make([]string, 0, len(algoValues))
for _, a := range algoValues {
algos = append(algos, strings.ToLower(string(a)))
}
return &cli.Command{
Name: "utils",
Usage: "utility helper CLI tool",
@@ -53,6 +73,69 @@ func UtilsCommand() *cli.Command {
Usage: "Convert legacy X-Amz-Meta.* xattrs into user.metadata JSON and remove legacy keys.",
Action: convertXattrMetadata,
},
{
Name: "convert-posix-dataset",
Aliases: []string{"cpd"},
Usage: "Make a preexisting posix dataset, not created through the gateway, usable by the posix backend.",
Description: `Walks the posix root directory and stores the metadata the posix backend
expects on buckets and, optionally, on objects.
Every top level directory is treated as a bucket: a private bucket ACL owned by
the --access-key-id account and the BucketOwnerEnforced object ownership are
stored on it. Top level files are ignored.
With --calculate-etag and/or --checksum-algorithm every regular file within a
bucket is read once to compute the MD5 ETag and/or the FULL_OBJECT checksum of
the requested algorithm.
Existing configuration is never overwritten: a bucket attribute, object ETag or
object checksum that is already present is left as is.
Metadata is stored in extended attributes, unless --sidecar-dir is set.
WARNING: for large datasets this command can take a really long time, as it has
to visit every file in the dataset. With --calculate-etag or --checksum-algorithm
every byte of every object is read and hashed. Tune --concurrency and
--read-buffer-size for the host running the conversion.`,
Action: convertPosixDataset,
Flags: []cli.Flag{
&cli.StringFlag{
Name: "posix-root-dir",
Usage: "the posix gateway root directory holding the dataset",
Required: true,
},
&cli.StringFlag{
Name: "sidecar-dir",
Usage: "store metadata in this sidecar directory instead of extended attributes; must be the same directory the gateway is run with",
},
&cli.StringFlag{
Name: "access-key-id",
Usage: "access key id of the account that owns the converted buckets",
DefaultText: "the root user access key",
},
&cli.BoolFlag{
Name: "calculate-etag",
Usage: "calculate and store the MD5 ETag of objects without one (reads every object; can take a long time on large datasets)",
},
&cli.StringFlag{
Name: "checksum-algorithm",
Usage: fmt.Sprintf("calculate and store the FULL_OBJECT checksum of objects without one, one of: %s (reads every object; can take a long time on large datasets)",
strings.Join(algos, ", ")),
},
&cli.IntFlag{
Name: "concurrency",
Usage: "maximum number of objects converted concurrently",
Value: runtime.NumCPU(),
DefaultText: "number of CPUs",
},
&cli.IntFlag{
Name: "read-buffer-size",
Usage: "per file read buffer size in bytes used for ETag and checksum calculation; peak buffer memory is concurrency * read-buffer-size",
Value: datasetDefaultReadBufferSize,
DefaultText: "1048576",
},
},
},
},
}
}
@@ -226,3 +309,443 @@ func convertXattrMetadata(ctx *cli.Context) error {
return nil
}
// Metadata attribute names and reserved directories
const (
datasetACLKey = "acl"
datasetOwnershipKey = "ownership"
datasetEtagKey = "etag"
datasetChecksumsKey = "checksums"
datasetMetaTmpDir = ".sgwtmp"
datasetObjLockDir = ".vgwlocks"
)
const (
// datasetLargeFileSize is the object size at and above which a warning
// is printed before hashing, since reading the file can take a while.
datasetLargeFileSize = 1 << 30 // 1 GiB
// datasetDefaultReadBufferSize is the default per-file read buffer size.
datasetDefaultReadBufferSize = 1 << 20 // 1 MiB
)
// datasetConverter stores the posix backend metadata on a preexisting
// dataset.
type datasetConverter struct {
root string
ms meta.MetadataStorer
acl []byte
calcEtag bool
checksumAlgo types.ChecksumAlgorithm
bufSize int
concurrency int
bucketsUpdated atomic.Int64
bucketsSkipped atomic.Int64
objectsUpdated atomic.Int64
objectsSkipped atomic.Int64
bytesHashed atomic.Int64
errors atomic.Int64
}
func convertPosixDataset(ctx *cli.Context) error {
root, err := filepath.Abs(ctx.String("posix-root-dir"))
if err != nil {
return fmt.Errorf("resolve posix root directory: %w", err)
}
fi, err := os.Stat(root)
if err != nil {
return fmt.Errorf("stat posix root directory: %w", err)
}
if !fi.IsDir() {
return fmt.Errorf("posix root %q is not a directory", root)
}
access := ctx.String("access-key-id")
if access == "" {
access = RootUserAccess
}
if access == "" {
return errors.New("bucket owner access key id is not set: provide --access-key-id or the root user access key")
}
concurrency := ctx.Int("concurrency")
if concurrency <= 0 {
return fmt.Errorf("concurrency must be positive, got %d", concurrency)
}
bufSize := ctx.Int("read-buffer-size")
if bufSize <= 0 {
return fmt.Errorf("read buffer size must be positive, got %d", bufSize)
}
var checksumAlgo types.ChecksumAlgorithm
if a := ctx.String("checksum-algorithm"); a != "" {
for _, algo := range types.ChecksumAlgorithmCrc32.Values() {
if strings.EqualFold(a, string(algo)) {
checksumAlgo = algo
break
}
}
if checksumAlgo == "" {
return fmt.Errorf("invalid checksum algorithm %q", a)
}
}
var ms meta.MetadataStorer
if sidecarDir := ctx.String("sidecar-dir"); sidecarDir != "" {
sidecarDir, err = filepath.Abs(sidecarDir)
if err != nil {
return fmt.Errorf("resolve sidecar directory: %w", err)
}
rel, err := filepath.Rel(root, sidecarDir)
if err == nil && rel != ".." && !strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return errors.New("sidecar directory cannot be inside the posix root directory")
}
sc, err := meta.NewSideCar(sidecarDir)
if err != nil {
return fmt.Errorf("init sidecar metadata: %w", err)
}
ms = sc
} else {
xm := meta.XattrMeta{}
if err := xm.Test(root); err != nil {
return fmt.Errorf("xattr check failed: %w", err)
}
ms = xm.WithRootDir(root)
}
acl, err := auth.UpdateACL(&auth.PutBucketAclInput{
ACL: types.BucketCannedACLPrivate,
}, auth.ACL{Owner: access}, nil)
if err != nil {
return fmt.Errorf("build bucket acl: %w", err)
}
utils.SetBucketNameValidationStrict(!DisableStrictBucketNames)
c := &datasetConverter{
root: root,
ms: ms,
acl: acl,
calcEtag: ctx.Bool("calculate-etag"),
checksumAlgo: checksumAlgo,
bufSize: bufSize,
concurrency: concurrency,
}
err = c.run(ctx.Context)
status := "is finished"
if errors.Is(err, context.Canceled) {
status = "was interrupted"
} else if err != nil {
status = "failed"
}
fmt.Printf("posix dataset conversion %s:\n"+
" directory: %s\n"+
" buckets updated: %d\n"+
" buckets skipped: %d\n"+
" objects updated: %d\n"+
" objects skipped: %d\n"+
" bytes hashed: %d\n"+
" errors: %d\n",
status, root, c.bucketsUpdated.Load(), c.bucketsSkipped.Load(),
c.objectsUpdated.Load(), c.objectsSkipped.Load(),
c.bytesHashed.Load(), c.errors.Load())
if err != nil {
return err
}
if n := c.errors.Load(); n > 0 {
return cli.Exit(fmt.Sprintf("posix dataset conversion finished with %d errors", n), 1)
}
return nil
}
func (c *datasetConverter) run(ctx context.Context) error {
ents, err := os.ReadDir(c.root)
if err != nil {
return fmt.Errorf("read posix root directory: %w", err)
}
hashObjects := c.calcEtag || c.checksumAlgo != ""
// jobs carries object paths relative to the root: <bucket>/<object>.
jobs := make(chan string, c.concurrency)
var wg sync.WaitGroup
if hashObjects {
for range c.concurrency {
wg.Go(func() {
// Each worker owns a single buffer, bounding the read
// buffer memory to concurrency * read-buffer-size.
var buf []byte
for job := range jobs {
if ctx.Err() != nil {
// Drain the queue without converting.
continue
}
if buf == nil {
buf = make([]byte, c.bufSize)
}
bucket, object, _ := strings.Cut(job, string(filepath.Separator))
err := c.convertObject(ctx, bucket, object, buf)
if errors.Is(err, context.Canceled) {
continue
}
if err != nil {
c.errorf("object %s: %v", job, err)
}
}
})
}
}
for _, ent := range ents {
if ctx.Err() != nil {
break
}
// Only real directories are buckets, root level files are ignored.
if !ent.IsDir() {
continue
}
bucket := ent.Name()
if bucket == datasetObjLockDir {
continue
}
if !utils.IsValidBucketName(bucket) {
fmt.Printf("skipping directory %q: invalid bucket name\n", bucket)
continue
}
if err := c.convertBucket(bucket); err != nil {
c.errorf("bucket %s: %v", bucket, err)
continue
}
if hashObjects {
fmt.Printf("converting objects of bucket %q\n", bucket)
c.walkBucket(ctx, bucket, jobs)
}
}
close(jobs)
wg.Wait()
return ctx.Err()
}
// convertBucket stores the bucket ACL and object ownership, unless already
// present.
func (c *datasetConverter) convertBucket(bucket string) error {
updated := false
for _, attr := range []struct {
key string
value []byte
}{
{datasetACLKey, c.acl},
{datasetOwnershipKey, []byte(types.ObjectOwnershipBucketOwnerEnforced)},
} {
exists, err := c.hasAttribute(bucket, "", attr.key)
if err != nil {
return err
}
if exists {
continue
}
if err := c.ms.StoreAttribute(nil, bucket, "", attr.key, attr.value); err != nil {
return fmt.Errorf("store %s: %w", attr.key, err)
}
updated = true
}
if updated {
c.bucketsUpdated.Add(1)
} else {
c.bucketsSkipped.Add(1)
}
return nil
}
// walkBucket queues every regular file in the bucket for conversion.
func (c *datasetConverter) walkBucket(ctx context.Context, bucket string, jobs chan<- string) {
bucketPath := filepath.Join(c.root, bucket)
err := filepath.WalkDir(bucketPath, func(path string, d fs.DirEntry, err error) error {
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil {
c.errorf("walk %s: %v", path, err)
return nil
}
if d.IsDir() {
if d.Name() == datasetMetaTmpDir && filepath.Dir(path) == bucketPath {
return fs.SkipDir
}
return nil
}
// Symlinks, sockets, devices, etc. are not objects.
if !d.Type().IsRegular() {
return nil
}
object, err := filepath.Rel(c.root, path)
if err != nil {
c.errorf("walk %s: %v", path, err)
return nil
}
select {
case jobs <- object:
case <-ctx.Done():
return ctx.Err()
}
return nil
})
if err != nil && !errors.Is(err, context.Canceled) {
c.errorf("walk bucket %s: %v", bucket, err)
}
}
// convertObject calculates and stores the missing ETag and checksum of an
// object, reading the file once for both.
func (c *datasetConverter) convertObject(ctx context.Context, bucket, object string, buf []byte) error {
needEtag := false
if c.calcEtag {
exists, err := c.hasAttribute(bucket, object, datasetEtagKey)
if err != nil {
return err
}
needEtag = !exists
}
needChecksum := false
if c.checksumAlgo != "" {
exists, err := c.hasAttribute(bucket, object, datasetChecksumsKey)
if err != nil {
return err
}
needChecksum = !exists
}
if !needEtag && !needChecksum {
c.objectsSkipped.Add(1)
return nil
}
f, err := os.Open(filepath.Join(c.root, bucket, object))
if err != nil {
return err
}
defer f.Close()
before, err := f.Stat()
if err != nil {
return err
}
if !before.Mode().IsRegular() {
c.objectsSkipped.Add(1)
return nil
}
if before.Size() >= datasetLargeFileSize {
fmt.Printf("warning: object %s/%s is %d bytes, calculating its ETag/checksum may take a while\n",
bucket, object, before.Size())
}
var md5Hash, checksumHash hash.Hash
var dst []io.Writer
if needEtag {
md5Hash = md5.New()
dst = append(dst, md5Hash)
}
if needChecksum {
checksumHash, err = utils.NewHash(utils.HashType(strings.ToLower(string(c.checksumAlgo))))
if err != nil {
return err
}
dst = append(dst, checksumHash)
}
n, err := io.CopyBuffer(io.MultiWriter(dst...), ctxReader{ctx: ctx, r: f}, buf)
if errors.Is(err, context.Canceled) {
return err
}
if err != nil {
return fmt.Errorf("read object data: %w", err)
}
c.bytesHashed.Add(n)
after, err := f.Stat()
if err != nil {
return err
}
if n != after.Size() || after.Size() != before.Size() || !after.ModTime().Equal(before.ModTime()) {
return errors.New("object was modified during conversion, skipping")
}
// The open file is only used by xattr storage, which does not need
// it to be writable.
if needChecksum {
sum := utils.Base64SumString(checksumHash.Sum(nil))
checksum := s3response.Checksum{
Algorithm: c.checksumAlgo,
Type: types.ChecksumTypeFullObject,
}
checksum.SetSum(c.checksumAlgo, &sum)
value, err := json.Marshal(checksum)
if err != nil {
return fmt.Errorf("marshal checksum: %w", err)
}
if err := c.ms.StoreAttribute(f, bucket, object, datasetChecksumsKey, value); err != nil {
return fmt.Errorf("store checksum: %w", err)
}
}
if needEtag {
etag := backend.GenerateEtag(md5Hash)
if err := c.ms.StoreAttribute(f, bucket, object, datasetEtagKey, []byte(etag)); err != nil {
return fmt.Errorf("store etag: %w", err)
}
}
c.objectsUpdated.Add(1)
return nil
}
func (c *datasetConverter) hasAttribute(bucket, object, attr string) (bool, error) {
_, err := c.ms.RetrieveAttribute(nil, bucket, object, attr)
if errors.Is(err, meta.ErrNoSuchKey) {
return false, nil
}
if err != nil {
return false, fmt.Errorf("get %s: %w", attr, err)
}
return true, nil
}
// ctxReader stops reading once ctx is canceled, so an interrupt does not
// wait for a large file to be fully hashed. Wrapping the file also hides
// *os.File's WriterTo, so io.CopyBuffer streams through the given buffer.
type ctxReader struct {
ctx context.Context
r io.Reader
}
func (cr ctxReader) Read(p []byte) (int, error) {
if err := cr.ctx.Err(); err != nil {
return 0, err
}
return cr.r.Read(p)
}
func (c *datasetConverter) errorf(format string, args ...any) {
c.errors.Add(1)
// Metadata storers report some failures (e.g. no space left on device)
// as S3 API errors, whose Error() is an XML document; print the
// description instead.
for i, arg := range args {
err, ok := arg.(error)
if !ok {
continue
}
if s3Err, ok := errors.AsType[s3err.S3Error](err); ok {
args[i] = strings.Replace(err.Error(), s3Err.Error(), s3Err.BaseError().Description, 1)
}
}
fmt.Fprintf(os.Stderr, "error: "+format+"\n", args...)
}
+33 -28
View File
@@ -87,34 +87,9 @@ var (
// checksum. If the provided sum is "", then the Sum() method can still
// be used to get the current checksum for the data read so far.
func NewHashReader(r io.Reader, expectedSum string, ht HashType) (*HashReader, error) {
var hash hash.Hash
switch ht {
case HashTypeContentMD5, HashTypeMd5:
hash = md5.New()
case HashTypeSha256Hex:
hash = sha256.New()
case HashTypeSha256:
hash = sha256.New()
case HashTypeSha1:
hash = sha1.New()
case HashTypeSha512:
hash = sha512.New()
case HashTypeCRC32:
hash = crc32.NewIEEE()
case HashTypeCRC32C:
hash = crc32.New(crc32.MakeTable(crc32.Castagnoli))
case HashTypeCRC64NVME:
hash = crc64.New(crc64NVMETable)
case HashTypeXXHASH64:
hash = xxhash.New()
case HashTypeXXHASH3:
hash = xxh3.New()
case HashTypeXXHASH128:
hash = xxh3.New128()
case HashTypeNone:
hash = noop{}
default:
return nil, errInvalidHashType
hash, err := NewHash(ht)
if err != nil {
return nil, err
}
return &HashReader{
@@ -125,6 +100,36 @@ func NewHashReader(r io.Reader, expectedSum string, ht HashType) (*HashReader, e
}, nil
}
// NewHash returns the hash.Hash implementing the given checksum algorithm.
func NewHash(ht HashType) (hash.Hash, error) {
switch ht {
case HashTypeContentMD5, HashTypeMd5:
return md5.New(), nil
case HashTypeSha256Hex, HashTypeSha256:
return sha256.New(), nil
case HashTypeSha1:
return sha1.New(), nil
case HashTypeSha512:
return sha512.New(), nil
case HashTypeCRC32:
return crc32.NewIEEE(), nil
case HashTypeCRC32C:
return crc32.New(crc32.MakeTable(crc32.Castagnoli)), nil
case HashTypeCRC64NVME:
return crc64.New(crc64NVMETable), nil
case HashTypeXXHASH64:
return xxhash.New(), nil
case HashTypeXXHASH3:
return xxh3.New(), nil
case HashTypeXXHASH128:
return xxh3.New128(), nil
case HashTypeNone:
return noop{}, nil
default:
return nil, errInvalidHashType
}
}
// Read allows *HashReader to be used as an io.Reader
func (hr *HashReader) Read(p []byte) (int, error) {
n, readerr := hr.r.Read(p)
+31
View File
@@ -770,6 +770,37 @@ type Checksum struct {
XXHASH128 *string
}
// SetSum stores sum in the field matching algo. A nil sum or an unknown
// algorithm leaves the checksum unchanged.
func (c *Checksum) SetSum(algo types.ChecksumAlgorithm, sum *string) {
if sum == nil {
return
}
switch algo {
case types.ChecksumAlgorithmCrc32:
c.CRC32 = sum
case types.ChecksumAlgorithmCrc32c:
c.CRC32C = sum
case types.ChecksumAlgorithmSha1:
c.SHA1 = sum
case types.ChecksumAlgorithmSha256:
c.SHA256 = sum
case types.ChecksumAlgorithmCrc64nvme:
c.CRC64NVME = sum
case types.ChecksumAlgorithmSha512:
c.SHA512 = sum
case types.ChecksumAlgorithmMd5:
c.MD5 = sum
case types.ChecksumAlgorithmXxhash64:
c.XXHASH64 = sum
case types.ChecksumAlgorithmXxhash3:
c.XXHASH3 = sum
case types.ChecksumAlgorithmXxhash128:
c.XXHASH128 = sum
}
}
// LocationConstraint represents the GetBucketLocation response
type LocationConstraint struct {
XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ LocationConstraint"`