mirror of
https://github.com/rclone/rclone.git
synced 2026-10-06 04:57:04 -04:00
operations: fix hang and lost context when moving a directory file by file
When the backend has no DirMove method, operations.DirMove moves the objects one by one. If one of the moves failed the movers stopped without receiving the rest of the objects, so DirMove blocked forever once the channel filled up. This could be seen renaming a directory on a mount of a backend without DirMove, such as s3. The moves were also done with context.Background() instead of the caller's context, so they didn't have the caller's config, couldn't be cancelled and weren't counted in the caller's stats.
This commit is contained in:
1 parent
2db37b14c2
commit
f48f2094a4
2 files changed
+87
-2
No files matched your search
@@ -1596,7 +1596,7 @@ func Rmdirs(ctx context.Context, f fs.Fs, dir string, leaveRoot bool) error {
|
||||
}
|
||||
fs.Debugf(nil, "removing %d level %d directories", len(dirs), level)
|
||||
sort.Strings(dirs)
|
||||
g, gCtx := errgroup.WithContext(context.Background())
|
||||
g, gCtx := errgroup.WithContext(ctx)
|
||||
g.SetLimit(ci.Checkers)
|
||||
for _, dir := range dirs {
|
||||
// End early if error
|
||||
@@ -2516,11 +2516,17 @@ func DirMove(ctx context.Context, f fs.Fs, srcRemote, dstRemote string) (err err
|
||||
return nil
|
||||
})
|
||||
}
|
||||
sending:
|
||||
for dir, entries := range tree {
|
||||
dstPath := dstRemote + dir[len(srcRemote):]
|
||||
for _, entry := range entries {
|
||||
if o, ok := entry.(fs.Object); ok {
|
||||
renames <- rename{o, path.Join(dstPath, path.Base(o.Remote()))}
|
||||
select {
|
||||
case renames <- rename{o, path.Join(dstPath, path.Base(o.Remote()))}:
|
||||
case <-gCtx.Done():
|
||||
// The workers have stopped so stop sending
|
||||
break sending
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1405,6 +1405,85 @@ func TestListFormat(t *testing.T) {
|
||||
|
||||
}
|
||||
|
||||
// noDirMoveFs wraps an Fs removing DirMove so operations.DirMove has
|
||||
// to move the objects one by one, and allows Move to be made to fail.
|
||||
type noDirMoveFs struct {
|
||||
fs.Fs
|
||||
features *fs.Features
|
||||
moveCtxs chan context.Context
|
||||
}
|
||||
|
||||
func (f *noDirMoveFs) Features() *fs.Features { return f.features }
|
||||
|
||||
func newNoDirMoveFs(t *testing.T, wrapped fs.Fs, failOn string) *noDirMoveFs {
|
||||
move := wrapped.Features().Move
|
||||
require.NotNil(t, move, "the test needs a backend with Move")
|
||||
f := &noDirMoveFs{
|
||||
Fs: wrapped,
|
||||
moveCtxs: make(chan context.Context, 100),
|
||||
}
|
||||
features := *wrapped.Features()
|
||||
features.DirMove = nil
|
||||
features.Move = func(ctx context.Context, src fs.Object, remote string) (fs.Object, error) {
|
||||
f.moveCtxs <- ctx
|
||||
if src.Remote() == failOn {
|
||||
return nil, errors.New("boom")
|
||||
}
|
||||
return move(ctx, src, remote)
|
||||
}
|
||||
f.features = &features
|
||||
return f
|
||||
}
|
||||
|
||||
// Check DirMove doesn't hang when moving the objects one by one and
|
||||
// one of the moves fails
|
||||
func TestDirMoveMoveError(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ctx, ci := fs.AddConfig(ctx)
|
||||
ci.Checkers = 2
|
||||
r := fstest.NewRun(t)
|
||||
// More files than the renames channel can hold
|
||||
for i := range 20 {
|
||||
r.WriteObject(ctx, fmt.Sprintf("A/file%d", i), "x", t1)
|
||||
}
|
||||
f := newNoDirMoveFs(t, r.Fremote, "A/file0")
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
done <- operations.DirMove(ctx, f, "A", "B")
|
||||
}()
|
||||
select {
|
||||
case err := <-done:
|
||||
require.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "boom")
|
||||
case <-time.After(30 * time.Second):
|
||||
t.Fatal("DirMove didn't return - deadlocked sending to the movers")
|
||||
}
|
||||
}
|
||||
|
||||
// Check the objects moved one by one by DirMove are moved with the
|
||||
// caller's context
|
||||
func TestDirMoveContext(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ctx, ci := fs.AddConfig(ctx)
|
||||
ci.Checkers = 2
|
||||
r := fstest.NewRun(t)
|
||||
r.WriteObject(ctx, "A/one", "one", t1)
|
||||
r.WriteObject(ctx, "A/two", "two", t1)
|
||||
f := newNoDirMoveFs(t, r.Fremote, "")
|
||||
|
||||
require.NoError(t, operations.DirMove(accounting.WithStatsGroup(ctx, "test-dirmove"), f, "A", "B"))
|
||||
close(f.moveCtxs)
|
||||
moves := 0
|
||||
for moveCtx := range f.moveCtxs {
|
||||
group, ok := accounting.StatsGroupFromContext(moveCtx)
|
||||
assert.True(t, ok, "move wasn't passed the caller's context")
|
||||
assert.Equal(t, "test-dirmove", group)
|
||||
moves++
|
||||
}
|
||||
assert.Equal(t, 2, moves)
|
||||
}
|
||||
|
||||
func TestDirMove(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
r := fstest.NewRun(t)
|
||||
|
||||
Reference in new issue
Block a user