Merge pull request #6185 from aduffeck/gather-space-info-in-parallel
Speed up me/drives by gathering space information in parallel
This commit is contained in:
@@ -34,6 +34,7 @@ import (
|
|||||||
settingsServiceExt "github.com/owncloud/ocis/v2/services/settings/pkg/store/defaults"
|
settingsServiceExt "github.com/owncloud/ocis/v2/services/settings/pkg/store/defaults"
|
||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
merrors "go-micro.dev/v4/errors"
|
merrors "go-micro.dev/v4/errors"
|
||||||
|
"golang.org/x/sync/errgroup"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -503,24 +504,70 @@ func (g Graph) UpdateDrive(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (g Graph) formatDrives(ctx context.Context, baseURL *url.URL, storageSpaces []*storageprovider.StorageSpace) ([]*libregraph.Drive, error) {
|
func (g Graph) formatDrives(ctx context.Context, baseURL *url.URL, storageSpaces []*storageprovider.StorageSpace) ([]*libregraph.Drive, error) {
|
||||||
responses := make([]*libregraph.Drive, 0, len(storageSpaces))
|
errg, ctx := errgroup.WithContext(ctx)
|
||||||
for _, storageSpace := range storageSpaces {
|
work := make(chan *storageprovider.StorageSpace, len(storageSpaces))
|
||||||
res, err := g.cs3StorageSpaceToDrive(ctx, baseURL, storageSpace)
|
results := make(chan *libregraph.Drive, len(storageSpaces))
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
// can't access disabled space
|
// Distribute work
|
||||||
if utils.ReadPlainFromOpaque(storageSpace.Opaque, "trashed") != "trashed" {
|
errg.Go(func() error {
|
||||||
res.Special = g.getExtendedSpaceProperties(ctx, baseURL, storageSpace)
|
defer close(work)
|
||||||
quota, err := g.getDriveQuota(ctx, storageSpace)
|
for _, space := range storageSpaces {
|
||||||
res.Quota = "a
|
select {
|
||||||
if err != nil {
|
case work <- space:
|
||||||
//logger.Debug().Err(err).Interface("id", sp.Id).Msg("error calling get quota on drive")
|
case <-ctx.Done():
|
||||||
return nil, err
|
return ctx.Err()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
responses = append(responses, res)
|
return nil
|
||||||
|
})
|
||||||
|
|
||||||
|
// Spawn workers that'll concurrently work the queue
|
||||||
|
numWorkers := 20
|
||||||
|
if len(storageSpaces) < numWorkers {
|
||||||
|
numWorkers = len(storageSpaces)
|
||||||
|
}
|
||||||
|
for i := 0; i < numWorkers; i++ {
|
||||||
|
errg.Go(func() error {
|
||||||
|
for storageSpace := range work {
|
||||||
|
res, err := g.cs3StorageSpaceToDrive(ctx, baseURL, storageSpace)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// can't access disabled space
|
||||||
|
if utils.ReadPlainFromOpaque(storageSpace.Opaque, "trashed") != "trashed" {
|
||||||
|
res.Special = g.getExtendedSpaceProperties(ctx, baseURL, storageSpace)
|
||||||
|
quota, err := g.getDriveQuota(ctx, storageSpace)
|
||||||
|
res.Quota = "a
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case results <- res:
|
||||||
|
case <-ctx.Done():
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// Wait for things to settle down, then close results chan
|
||||||
|
go func() {
|
||||||
|
_ = errg.Wait() // error is checked later
|
||||||
|
close(results)
|
||||||
|
}()
|
||||||
|
|
||||||
|
responses := make([]*libregraph.Drive, len(storageSpaces))
|
||||||
|
i := 0
|
||||||
|
for r := range results {
|
||||||
|
responses[i] = r
|
||||||
|
i++
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := errg.Wait(); err != nil {
|
||||||
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return responses, nil
|
return responses, nil
|
||||||
|
|||||||
Reference in New Issue
Block a user