diff --git a/backend/posix/posix.go b/backend/posix/posix.go index f34446cb..eef0537d 100644 --- a/backend/posix/posix.go +++ b/backend/posix/posix.go @@ -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) diff --git a/cmd/internal/gwcli/utils.go b/cmd/internal/gwcli/utils.go index 5f259d76..d249debd 100644 --- a/cmd/internal/gwcli/utils.go +++ b/cmd/internal/gwcli/utils.go @@ -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: /. + 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...) +} diff --git a/s3api/utils/csum-reader.go b/s3api/utils/csum-reader.go index 0bc6bc50..7eb7d131 100644 --- a/s3api/utils/csum-reader.go +++ b/s3api/utils/csum-reader.go @@ -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) diff --git a/s3response/s3response.go b/s3response/s3response.go index 357f8bc3..d9a24a17 100644 --- a/s3response/s3response.go +++ b/s3response/s3response.go @@ -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"`