diff --git a/cmd/versitygw/test.go b/cmd/versitygw/test.go index d12418fc..92b06b9f 100644 --- a/cmd/versitygw/test.go +++ b/cmd/versitygw/test.go @@ -8,9 +8,19 @@ import ( ) var ( - awsID string - awsSecret string - endpoint string + awsID string + awsSecret string + endpoint string + prefix string + dstBucket string + partSize int64 + objSize int64 + concurrency int + files int + upload bool + download bool + pathStyle bool + checksumDisable bool ) func testCommand() *cli.Command { @@ -153,6 +163,106 @@ func initTestCommands() []*cli.Command { Description: `Runs all the available tests to test the full flow of the gateway.`, Action: getAction(integration.TestFullFlow), }, + { + Name: "bench", + Usage: "Runs download/upload performance test on the gateway", + Description: `Uploads/downloads some number(specified by flags) of files with some capacity(bytes). + Logs the results to the console`, + Flags: []cli.Flag{ + &cli.IntFlag{ + Name: "files", + Usage: "Number of objects to read/write", + Value: 1, + Destination: &files, + }, + &cli.Int64Flag{ + Name: "objsize", + Usage: "Uploading object size", + Value: 0, + Destination: &objSize, + }, + &cli.StringFlag{ + Name: "prefix", + Usage: "Object name prefix", + Destination: &prefix, + }, + &cli.BoolFlag{ + Name: "upload", + Usage: "Upload data to the gateway", + Value: false, + Destination: &upload, + }, + &cli.BoolFlag{ + Name: "download", + Usage: "Download data to the gateway", + Value: false, + Destination: &download, + }, + &cli.StringFlag{ + Name: "bucket", + Usage: "Destination bucket name to read/write data", + Destination: &dstBucket, + }, + &cli.Int64Flag{ + Name: "partSize", + Usage: "Upload/download size per thread", + Value: 64 * 1024 * 1024, + Destination: &partSize, + }, + &cli.IntFlag{ + Name: "concurrency", + Usage: "Upload/download threads per object", + Value: 1, + Destination: &concurrency, + }, + &cli.BoolFlag{ + Name: "pathStyle", + Usage: "Use Pathstyle bucket addressing", + Value: false, + Destination: &pathStyle, + }, + &cli.BoolFlag{ + Name: "checksumDis", + Usage: "Disable server checksum", + Value: false, + Destination: &checksumDisable, + }, + }, + Action: func(ctx *cli.Context) error { + if upload && download { + return fmt.Errorf("must only specify one of upload or download") + } + if !upload && !download { + return fmt.Errorf("must specify one of upload or download") + } + + if dstBucket == "" { + return fmt.Errorf("must specify bucket") + } + + opts := []integration.Option{ + integration.WithAccess(awsID), + integration.WithSecret(awsSecret), + integration.WithRegion(region), + integration.WithEndpoint(endpoint), + integration.WithConcurrency(concurrency), + integration.WithPartSize(partSize), + } + if debug { + opts = append(opts, integration.WithDebug()) + } + if pathStyle { + opts = append(opts, integration.WithPathStyle()) + } + if checksumDisable { + opts = append(opts, integration.WithDisableChecksum()) + } + + s3conf := integration.NewS3Conf(opts...) + + return integration.TestPerformance(s3conf, upload, download, files, objSize, dstBucket, prefix) + }, + }, } } diff --git a/integration/reader.go b/integration/data-io.go similarity index 59% rename from integration/reader.go rename to integration/data-io.go index 7a8d2d27..491059fb 100644 --- a/integration/reader.go +++ b/integration/data-io.go @@ -38,6 +38,26 @@ func (r *RReader) Sum() []byte { return r.hash.Sum(nil) } +type ZReader struct { + buf []byte + dataleft int +} + +func NewZeroReader(totalsize, bufsize int) *ZReader { + b := make([]byte, bufsize) + return &ZReader{buf: b, dataleft: totalsize} +} + +func (r *ZReader) Read(p []byte) (int, error) { + n := min(len(p), len(r.buf), r.dataleft) + r.dataleft -= n + err := error(nil) + if n == 0 { + err = io.EOF + } + return copy(p, r.buf[:n]), err +} + func min(values ...int) int { if len(values) == 0 { return 0 @@ -52,3 +72,13 @@ func min(values ...int) int { return min } + +type NW struct{} + +func NewNullWriter() NW { + return NW{} +} + +func (NW) WriteAt(p []byte, off int64) (n int, err error) { + return len(p), nil +} diff --git a/integration/s3conf.go b/integration/s3conf.go index 2533d49e..89ea13f8 100644 --- a/integration/s3conf.go +++ b/integration/s3conf.go @@ -2,6 +2,7 @@ package integration import ( "context" + "io" "log" "net/http" "os" @@ -10,6 +11,8 @@ import ( v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4" "github.com/aws/aws-sdk-go-v2/config" "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/aws/smithy-go/middleware" ) @@ -26,10 +29,7 @@ type S3Conf struct { } func NewS3Conf(opts ...Option) *S3Conf { - s := &S3Conf{ - PartSize: 64 * 1024 * 1024, // 64B default chunksize - Concurrency: 1, // 1 default concurrency - } + s := &S3Conf{} for _, opt := range opts { opt(s) @@ -123,3 +123,31 @@ func (c *S3Conf) Config() aws.Config { return cfg } + +func (c *S3Conf) UploadData(r io.Reader, bucket, object string) error { + uploader := manager.NewUploader(s3.NewFromConfig(c.Config())) + uploader.PartSize = c.PartSize + uploader.Concurrency = c.Concurrency + + upinfo := &s3.PutObjectInput{ + Body: r, + Bucket: &bucket, + Key: &object, + } + + _, err := uploader.Upload(context.Background(), upinfo) + return err +} + +func (c *S3Conf) DownloadData(w io.WriterAt, bucket, object string) (int64, error) { + downloader := manager.NewDownloader(s3.NewFromConfig(c.Config())) + downloader.PartSize = c.PartSize + downloader.Concurrency = c.Concurrency + + downinfo := &s3.GetObjectInput{ + Bucket: &bucket, + Key: &object, + } + + return downloader.Download(context.Background(), w, downinfo) +} diff --git a/integration/tests.go b/integration/tests.go index f981a3a6..d12347d8 100644 --- a/integration/tests.go +++ b/integration/tests.go @@ -7,11 +7,12 @@ import ( "crypto/sha256" "fmt" "io" + "math" "os" "strings" + "sync" "time" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/aws/aws-sdk-go-v2/service/s3/types" ) @@ -193,7 +194,7 @@ func TestPutGetMPObject(s *S3Conf) { dr := NewDataReader(datalen, 5*1024*1024) WithPartSize(5 * 1024 * 1024) s.PartSize = 5 * 1024 * 1024 - err = uploadData(s, dr, bucket, name) + err = s.UploadData(dr, bucket, name) if err != nil { failF("%v: %v", testname, err) return @@ -258,21 +259,6 @@ func isEqual(a, b []byte) bool { return true } -func uploadData(s *S3Conf, r io.Reader, bucket, object string) error { - uploader := manager.NewUploader(s3.NewFromConfig(s.Config())) - uploader.PartSize = s.PartSize - uploader.Concurrency = s.Concurrency - - upinfo := &s3.PutObjectInput{ - Body: r, - Bucket: &bucket, - Key: &object, - } - - _, err := uploader.Upload(context.Background(), upinfo) - return err -} - func TestPutDirObject(s *S3Conf) { testname := "test put directory object" runF(testname) @@ -1148,6 +1134,77 @@ func TestInvalidMultiParts(s *S3Conf) { passF(testname) } +type prefResult struct { + elapsed time.Duration + size int64 + err error +} + +func TestPerformance(s *S3Conf, upload, download bool, files int, objectSize int64, bucket, prefix string) error { + var sg sync.WaitGroup + results := make([]prefResult, files) + start := time.Now() + if upload { + if objectSize == 0 { + return fmt.Errorf("must specify object size for upload") + } + + if objectSize > (int64(10000) * s.PartSize) { + return fmt.Errorf("object size can not exceed 10000 * chunksize") + } + + runF("performance test: upload/download objects") + + for i := 0; i < files; i++ { + sg.Add(1) + go func(i int) { + var r io.Reader = NewDataReader(int(objectSize), int(s.PartSize)) + + start := time.Now() + err := s.UploadData(r, bucket, fmt.Sprintf("%v%v", prefix, i)) + results[i].elapsed = time.Since(start) + results[i].err = err + results[i].size = objectSize + sg.Done() + }(i) + } + } + if download { + for i := 0; i < files; i++ { + sg.Add(1) + go func(i int) { + nw := NewNullWriter() + start := time.Now() + n, err := s.DownloadData(nw, bucket, fmt.Sprintf("%v%v", prefix, i)) + results[i].elapsed = time.Since(start) + results[i].err = err + results[i].size = n + sg.Done() + }(i) + } + } + sg.Wait() + elapsed := time.Since(start) + + var tot int64 + for i, res := range results { + if res.err != nil { + failF("%v: %v\n", i, res.err) + break + } + tot += res.size + fmt.Printf("%v: %v in %v (%v MB/s)\n", + i, res.size, res.elapsed, + int(math.Ceil(float64(res.size)/res.elapsed.Seconds())/1048576)) + } + + fmt.Println() + passF("run perf: %v in %v (%v MB/s)\n", + tot, elapsed, int(math.Ceil(float64(tot)/elapsed.Seconds())/1048576)) + + return nil +} + // Full flow test func TestFullFlow(s *S3Conf) { // TODO: add more test cases to get 100% coverage