diff --git a/backend/hidrive/helpers.go b/backend/hidrive/helpers.go index e7afc8a6d..c64d53171 100644 --- a/backend/hidrive/helpers.go +++ b/backend/hidrive/helpers.go @@ -24,6 +24,7 @@ import ( "github.com/rclone/rclone/fs" "github.com/rclone/rclone/fs/accounting" "github.com/rclone/rclone/fs/fserrors" + "github.com/rclone/rclone/lib/pool" "github.com/rclone/rclone/lib/ranges" "github.com/rclone/rclone/lib/readers" "github.com/rclone/rclone/lib/rest" @@ -454,17 +455,17 @@ func (f *Fs) deleteObject(ctx context.Context, path string) error { } // createFile creates a file at the given path -// with the content of the io.ReadSeeker. +// with the content of the buffer. // This guarantees that existing files will not be overwritten. // The maximum size of the content is limited by MaximumUploadBytes. -// The io.ReadSeeker should be resettable by seeking to its start. +// The caller remains responsible for closing the buffer. // If modTime is not the zero time instant, // it will be set as the file's modification time after the operation. // // This returns fs.ErrorDirNotFound // if the parent directory of the file is not found. // This returns ErrorFileExists if a file already exists at the specified path. -func (f *Fs) createFile(ctx context.Context, path string, content io.ReadSeeker, modTime time.Time, onExist OnExistAction) (*api.HiDriveObject, error) { +func (f *Fs) createFile(ctx context.Context, path string, content *pool.RW, modTime time.Time, onExist OnExistAction) (*api.HiDriveObject, error) { parameters := api.NewQueryParameters() parameters.SetFileInDirectory(path) if onExist == AutoNameOnExist { @@ -479,12 +480,14 @@ func (f *Fs) createFile(ctx context.Context, path string, content io.ReadSeeker, } } + contentLength := content.Size() opts := rest.Opts{ - Method: "POST", - Path: "/file", - Body: content, - ContentType: "application/octet-stream", - Parameters: parameters.Values, + Method: "POST", + Path: "/file", + Body: content, + ContentType: "application/octet-stream", + ContentLength: &contentLength, + Parameters: parameters.Values, } var result api.HiDriveObject @@ -839,6 +842,17 @@ func createHiDriveScopes(role string, access string) []string { return []string{} } +// unwrapAccounting splits any transfer accounting off reader, +// returning the raw stream and the accounting (nil if there is none), +// so the stream can be buffered and accounted when the buffer is uploaded. +func unwrapAccounting(reader io.Reader) (unwrapped io.Reader, acc *accounting.Account) { + // Any kind of accounter re-wraps into the one type UnWrapAccounting + // knows how to take the *Account out of. + unwrapped, wrap := accounting.UnWrap(reader) + _, acc = accounting.UnWrapAccounting(wrap(unwrapped)) + return unwrapped, acc +} + // accountedReadSeeker reads through any accounting wrapped around a // buffered reader while seeking the buffer underneath it, // so a retry can rewind the buffer without copying it. diff --git a/backend/hidrive/hidrive.go b/backend/hidrive/hidrive.go index cb512bcb0..a84dbaa09 100644 --- a/backend/hidrive/hidrive.go +++ b/backend/hidrive/hidrive.go @@ -28,6 +28,7 @@ import ( "github.com/rclone/rclone/fs/config/obscure" "github.com/rclone/rclone/fs/fserrors" "github.com/rclone/rclone/fs/hash" + "github.com/rclone/rclone/lib/multipart" "github.com/rclone/rclone/lib/oauthutil" "github.com/rclone/rclone/lib/pacer" "github.com/rclone/rclone/lib/rest" @@ -514,9 +515,24 @@ func (f *Fs) PutUnchecked(ctx context.Context, in io.Reader, src fs.ObjectInfo, // (i.e. everything up to the cutoff) in the first request, // avoids files being created on upload failure for small files. // (As opposed to creating an empty file and then uploading the content.) - tmpReader, bytesRead, err := readerForChunk(in, int(f.opt.UploadCutoff)) - cutoffReader := cachedReader(tmpReader) + // + // Only buffer as much as the source declares it has, + // so small files do not cost the whole cutoff in memory. + prefixSize := int64(f.opt.UploadCutoff) + if size := src.Size(); size >= 0 && size < prefixSize { + prefixSize = size + } + unwrapped, acc := unwrapAccounting(in) + cutoffReader := multipart.NewRW() + if acc != nil { + cutoffReader.SetAccounting(acc.AccountRead) + } + bytesRead, err := io.CopyN(cutoffReader, unwrapped, prefixSize) + if err == io.EOF { + err = nil + } if err != nil { + _ = cutoffReader.Close() return nil, err } @@ -540,6 +556,7 @@ func (f *Fs) PutUnchecked(ctx context.Context, in io.Reader, src fs.ObjectInfo, } return false, createErr }) + _ = cutoffReader.Close() if err != nil { return nil, err @@ -556,7 +573,7 @@ func (f *Fs) PutUnchecked(ctx context.Context, in io.Reader, src fs.ObjectInfo, } // If there is more left to write, o.Update needs to skip ahead. // Use a fs.SeekOption with the current offset to do this. - options = append(options, &fs.SeekOption{Offset: int64(bytesRead)}) + options = append(options, &fs.SeekOption{Offset: bytesRead}) err = o.Update(ctx, in, src, options...) if err == nil {