mirror of
https://github.com/rclone/rclone.git
synced 2026-10-09 14:39:28 -04:00
serve ftp: fix transfers failing after 5 minutes with --auth-proxy
With --auth-proxy the backend of each user is shut down once it has
been unused for 5 minutes. Only the start of each FTP command counted
as a use, so an upload or download lasting longer than that had its
backend shut down under it and failed with "context canceled".
Each FTP command now holds the backend it uses until it has finished,
which for a download is when the file is closed.
This was introduced in v1.75.1 by
f425f8d46 serve: refactor VFS and proxy handling into Provider
This commit is contained in:
1 parent
408cc55171
commit
ee0a70f875
2 files changed
+147
-6
No files matched your search
+48
-6
@@ -16,6 +16,7 @@ import (
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/rclone/rclone/cmd"
|
||||
@@ -358,13 +359,27 @@ func (d *driver) CheckPasswd(sctx *ftp.Context, user, pass string) (ok bool, err
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// getVFS returns the VFS for this connection.
|
||||
// getVFS returns the VFS for this connection, held so it can't be
|
||||
// shut down while it is in use. The caller must call Shutdown on it
|
||||
// when it has finished with it.
|
||||
//
|
||||
// In proxy mode, getVFS calls proxy.Call on each FTP command which refreshes
|
||||
// the proxy cache timer (like http/webdav). Therefore, connection-level pinning
|
||||
// is not used; only individual transfers exceeding the cache expiry window
|
||||
// could be affected.
|
||||
// In proxy mode, getVFS calls proxy.Call on each FTP command which
|
||||
// refreshes the proxy cache timer (like http/webdav). The proxy shuts
|
||||
// the VFS down when that timer expires, which a single transfer can
|
||||
// outlast, so the timer alone isn't enough to keep the VFS alive.
|
||||
func (d *driver) getVFS(sctx *ftp.Context) (VFS *vfs.VFS, err error) {
|
||||
VFS, err = d.findVFS(sctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !VFS.Hold() {
|
||||
return nil, errors.New("VFS has been shut down")
|
||||
}
|
||||
return VFS, nil
|
||||
}
|
||||
|
||||
// findVFS returns the VFS for this connection without holding it.
|
||||
func (d *driver) findVFS(sctx *ftp.Context) (VFS *vfs.VFS, err error) {
|
||||
if !d.provider.IsProxy() {
|
||||
// If no proxy always use the same VFS
|
||||
return d.provider.VFS(), nil
|
||||
@@ -392,6 +407,7 @@ func (d *driver) Stat(sctx *ftp.Context, path string) (fi iofs.FileInfo, err err
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
n, err := VFS.Stat(path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -406,6 +422,7 @@ func (d *driver) ChangeDir(sctx *ftp.Context, path string) (err error) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
n, err := VFS.Stat(path)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -423,6 +440,7 @@ func (d *driver) ListDir(sctx *ftp.Context, path string, callback func(iofs.File
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
node, err := VFS.Stat(path)
|
||||
if err == vfs.ENOENT {
|
||||
return errors.New("directory not found")
|
||||
@@ -461,6 +479,7 @@ func (d *driver) DeleteDir(sctx *ftp.Context, path string) (err error) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
node, err := VFS.Stat(path)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -482,6 +501,7 @@ func (d *driver) DeleteFile(sctx *ftp.Context, path string) (err error) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
node, err := VFS.Stat(path)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -503,6 +523,7 @@ func (d *driver) Rename(sctx *ftp.Context, oldName, newName string) (err error)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
return VFS.Rename(oldName, newName)
|
||||
}
|
||||
|
||||
@@ -513,6 +534,7 @@ func (d *driver) MakeDir(sctx *ftp.Context, path string) (err error) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
dir, leaf, err := VFS.StatParent(path)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -528,6 +550,12 @@ func (d *driver) GetFile(sctx *ftp.Context, path string, offset int64) (size int
|
||||
if err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
// The returned file takes over the hold on the VFS
|
||||
defer func() {
|
||||
if err != nil {
|
||||
VFS.Shutdown()
|
||||
}
|
||||
}()
|
||||
node, err := VFS.Stat(path)
|
||||
if err == vfs.ENOENT {
|
||||
fs.Infof(path, "File not found")
|
||||
@@ -552,7 +580,20 @@ func (d *driver) GetFile(sctx *ftp.Context, path string, offset int64) (size int
|
||||
tr := accounting.GlobalStats().NewTransferRemoteSize(path, node.Size(), d.f, nil)
|
||||
defer tr.Done(d.ctx, nil)
|
||||
|
||||
return node.Size(), handle, nil
|
||||
return node.Size(), &heldFile{ReadCloser: handle, VFS: VFS}, nil
|
||||
}
|
||||
|
||||
// heldFile is an open file which holds its VFS until it is closed.
|
||||
type heldFile struct {
|
||||
io.ReadCloser
|
||||
VFS *vfs.VFS
|
||||
release sync.Once
|
||||
}
|
||||
|
||||
// Close closes the file and releases the hold on the VFS.
|
||||
func (f *heldFile) Close() error {
|
||||
defer f.release.Do(f.VFS.Shutdown)
|
||||
return f.ReadCloser.Close()
|
||||
}
|
||||
|
||||
// PutFile upload a file
|
||||
@@ -564,6 +605,7 @@ func (d *driver) PutFile(sctx *ftp.Context, path string, data io.Reader, offset
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer VFS.Shutdown()
|
||||
fi, err := VFS.Stat(path)
|
||||
if err == nil {
|
||||
isExist = true
|
||||
|
||||
@@ -9,8 +9,14 @@ package ftp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
ftpclient "github.com/jlaffaye/ftp"
|
||||
_ "github.com/rclone/rclone/backend/local"
|
||||
"github.com/rclone/rclone/cmd/serve/proxy"
|
||||
"github.com/rclone/rclone/cmd/serve/servetest"
|
||||
@@ -19,6 +25,7 @@ import (
|
||||
"github.com/rclone/rclone/fs/config/obscure"
|
||||
"github.com/rclone/rclone/fs/rc"
|
||||
"github.com/rclone/rclone/lib/israce"
|
||||
"github.com/rclone/rclone/lib/random"
|
||||
"github.com/rclone/rclone/vfs"
|
||||
"github.com/rclone/rclone/vfs/vfscommon"
|
||||
"github.com/stretchr/testify/assert"
|
||||
@@ -149,3 +156,95 @@ func TestNewServerError(t *testing.T) {
|
||||
assert.Nil(t, d)
|
||||
assert.Equal(t, before, vfs.ActiveCount(), "VFS leaked after failed server creation")
|
||||
}
|
||||
|
||||
// TestAuthProxyTransferOutlivesCache checks transfers in progress
|
||||
// carry on working when the auth proxy drops their VFS from its cache,
|
||||
// as it does when a transfer takes longer than the cache expiry time.
|
||||
func TestAuthProxyTransferOutlivesCache(t *testing.T) {
|
||||
const addr = "127.0.0.1:" + testPORT
|
||||
root := t.TempDir()
|
||||
contents := random.String(32 * 1024 * 1024)
|
||||
require.NoError(t, os.WriteFile(filepath.Join(root, "download.bin"), []byte(contents), 0666))
|
||||
|
||||
prog, err := filepath.Abs("../servetest/proxy_code.go")
|
||||
require.NoError(t, err)
|
||||
opt := Opt
|
||||
opt.ListenAddr = addr
|
||||
opt.PassivePorts = testPASSIVEPORTRANGE
|
||||
proxyOpt := proxy.Opt
|
||||
proxyOpt.AuthProxy = "go run " + prog + " " + root
|
||||
d, err := newServer(context.Background(), nil, &opt, &vfscommon.Opt, &proxyOpt)
|
||||
require.NoError(t, err)
|
||||
quit := make(chan struct{})
|
||||
go func() {
|
||||
assert.NoError(t, d.Serve())
|
||||
close(quit)
|
||||
}()
|
||||
defer func() {
|
||||
assert.NoError(t, d.Shutdown())
|
||||
<-quit
|
||||
}()
|
||||
|
||||
var c *ftpclient.ServerConn
|
||||
require.Eventually(t, func() bool {
|
||||
c, err = ftpclient.Dial(addr)
|
||||
return err == nil
|
||||
}, 10*time.Second, 10*time.Millisecond)
|
||||
defer func() { _ = c.Quit() }()
|
||||
require.NoError(t, c.Login(testUSER, testPASS))
|
||||
|
||||
// Only the IP of the address is used, which is the same as the client's
|
||||
_, vfsKey, err := d.provider.Proxy().Call(testUSER, testPASS, false, addr)
|
||||
require.NoError(t, err)
|
||||
|
||||
// expire waits for a transfer to be using the VFS, then drops
|
||||
// everything from the proxy's cache as if it had expired.
|
||||
expire := func(t *testing.T) {
|
||||
var VFS *vfs.VFS
|
||||
require.Eventually(t, func() bool {
|
||||
VFS = d.provider.Proxy().Get(vfsKey)
|
||||
return VFS != nil && VFS.Stats()["inUse"] == int32(2)
|
||||
}, 10*time.Second, 10*time.Millisecond, "transfer isn't holding the VFS")
|
||||
d.provider.Proxy().Shutdown()
|
||||
assert.Equal(t, int32(1), VFS.Stats()["inUse"], "VFS not held by the transfer alone")
|
||||
}
|
||||
|
||||
t.Run("Download", func(t *testing.T) {
|
||||
resp, err := c.Retr("download.bin")
|
||||
require.NoError(t, err)
|
||||
defer func() { _ = resp.Close() }()
|
||||
start := make([]byte, 1024)
|
||||
_, err = io.ReadFull(resp, start)
|
||||
require.NoError(t, err)
|
||||
expire(t)
|
||||
rest, err := io.ReadAll(resp)
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, resp.Close())
|
||||
assert.True(t, contents == string(start)+string(rest), "download corrupted")
|
||||
})
|
||||
|
||||
t.Run("Upload", func(t *testing.T) {
|
||||
pr, pw := io.Pipe()
|
||||
var storErr error
|
||||
var wg sync.WaitGroup
|
||||
wg.Go(func() {
|
||||
storErr = c.Stor("upload.bin", pr)
|
||||
})
|
||||
// Finish the upload before the connection is used again
|
||||
defer func() {
|
||||
_ = pw.Close()
|
||||
wg.Wait()
|
||||
}()
|
||||
_, err := io.WriteString(pw, contents[:1024*1024])
|
||||
require.NoError(t, err)
|
||||
expire(t)
|
||||
_, err = io.WriteString(pw, contents[1024*1024:])
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, pw.Close())
|
||||
wg.Wait()
|
||||
require.NoError(t, storErr)
|
||||
got, err := os.ReadFile(filepath.Join(root, "upload.bin"))
|
||||
require.NoError(t, err)
|
||||
assert.True(t, contents == string(got), "upload corrupted")
|
||||
})
|
||||
}
|
||||
Reference in new issue
Block a user