2
2
mirror of https://github.com/octoleo/restic.git synced 2025-01-25 16:18:34 +00:00
restic/internal/repository/parallel.go

77 lines
1.7 KiB
Go
Raw Normal View History

2015-07-04 17:47:42 +02:00
package repository
import (
2017-06-04 11:16:55 +02:00
"context"
2015-07-04 17:47:42 +02:00
"sync"
2017-07-23 14:21:03 +02:00
"github.com/restic/restic/internal/debug"
2017-07-24 17:42:25 +02:00
"github.com/restic/restic/internal/restic"
2015-07-04 17:47:42 +02:00
)
// ParallelWorkFunc gets one file ID to work on. If an error is returned,
2017-06-04 11:16:55 +02:00
// processing stops. When the contect is cancelled the function should return.
type ParallelWorkFunc func(ctx context.Context, id string) error
2016-08-31 22:39:36 +02:00
// ParallelIDWorkFunc gets one restic.ID to work on. If an error is returned,
2017-06-04 11:16:55 +02:00
// processing stops. When the context is cancelled the function should return.
type ParallelIDWorkFunc func(ctx context.Context, id restic.ID) error
2015-07-04 17:47:42 +02:00
// FilesInParallel runs n workers of f in parallel, on the IDs that
// repo.List(t) yield. If f returns an error, the process is aborted and the
// first error is returned.
2017-06-04 11:16:55 +02:00
func FilesInParallel(ctx context.Context, repo restic.Lister, t restic.FileType, n uint, f ParallelWorkFunc) error {
2015-07-04 17:47:42 +02:00
wg := &sync.WaitGroup{}
2017-06-04 11:16:55 +02:00
ch := repo.List(ctx, t)
2015-07-04 17:47:42 +02:00
errors := make(chan error, n)
for i := 0; uint(i) < n; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case id, ok := <-ch:
2015-07-04 17:47:42 +02:00
if !ok {
return
}
2017-06-04 11:16:55 +02:00
err := f(ctx, id)
2015-07-04 17:47:42 +02:00
if err != nil {
errors <- err
return
}
2017-06-04 11:16:55 +02:00
case <-ctx.Done():
2015-07-04 17:47:42 +02:00
return
}
}
}()
}
wg.Wait()
select {
case err := <-errors:
return err
default:
break
}
return nil
}
2016-08-31 22:39:36 +02:00
// ParallelWorkFuncParseID converts a function that takes a restic.ID to a
// function that takes a string. Filenames that do not parse as a restic.ID
// are ignored.
func ParallelWorkFuncParseID(f ParallelIDWorkFunc) ParallelWorkFunc {
2017-06-04 11:16:55 +02:00
return func(ctx context.Context, s string) error {
2016-08-31 22:39:36 +02:00
id, err := restic.ParseID(s)
if err != nil {
2016-09-27 22:35:08 +02:00
debug.Log("invalid ID %q: %v", id, err)
return nil
}
2017-06-04 11:16:55 +02:00
return f(ctx, id)
}
}