fix: clean up a fragment left by a failed upload, and skip unwritable names

rclone writes straight to the final remote name on any backend that does not
advertise PartialUploads, and cleans up after a failed Put only when it did
not. Pikpak advertises neither, so a transfer that died halfway left a
fragment under exactly the name verification matches on — counted archived by
that run and every run after it, with the local copy already deleted. A failed
move now looks for that object and removes it, leaving a complete one alone
since pikpak's async commit can still land it correctly.

An unwritable filename ended the whole walk, so one hostile name could strand
every message behind it. It is skipped and reported instead. The comment there
had described that behaviour all along.

Also: remove the part file when promoting it fails, since the caller hands the
reservation back and the cap would stay over-committed; report upload failures
alongside a download error rather than instead of it, which on Ctrl-C hid that
finished files had been discarded.

Run had no test of its own because it called Download directly. That step is
now indirected, covering the properties only the composition has: uploads
closed after the last send, the budget balanced across failures, and a tripped
breaker halting downloads rather than walking the whole chat.
This commit is contained in:
tiennm99 committed 2026-09-06 19:51:24 +07:00
1 parent bd816156eb
commit 3e75138c59
7 files changed
+542 -58

No files matched your search

+11 -1
View File
@@ -103,7 +103,10 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo
if err == nil {
err = it.failure
}
return outcomes, stats, err
// Skipped items are reported alongside whatever else happened rather than
// instead of it: the run did real work, and the caller still needs to know
// these messages were never attempted.
return outcomes, stats, errors.Join(append([]error{err}, it.skipped...)...)
}
// finish closes a downloaded file and either promotes it or removes it.
@@ -143,6 +146,13 @@ func finish(staging string, e *elem, downloadErr error) error {
}
if err := os.Rename(part, finalPath(staging, e.item)); err != nil {
// The part file goes too. The caller treats this as a failure and hands
// the byte reservation back, so leaving the file on disk would put the
// staging cap permanently over-committed by its size.
if rerr := os.Remove(part); rerr != nil && !os.IsNotExist(rerr) {
return errors.Join(fmt.Errorf("promote %q: %w", part, err),
fmt.Errorf("and it is still in staging: %w", rerr))
}
return fmt.Errorf("promote %q: %w", part, err)
}
return nil
+47 -15
View File
@@ -129,9 +129,47 @@ func TestSweepPartialsOnMissingDirectory(t *testing.T) {
}
}
// An unwritable name must stop the iterator with a message naming the message,
// rather than surfacing as a bare os.Create failure later.
func TestElemIterRejectsUnsafeNames(t *testing.T) {
// An unwritable name must be refused with a message naming the message, rather
// than surfacing as a bare os.Create failure later — and it must be skipped, not
// treated as the end of the walk. One hostile filename cannot be allowed to
// strand every message behind it.
func TestElemIterSkipsUnsafeNamesAndKeepsGoing(t *testing.T) {
staging := t.TempDir()
bad := testItem(t, 7, "../../escape.conf", 10)
good := testItem(t, 8, "fine.mp4", 10)
seq := func(yield func(tgsource.Item, error) bool) {
if !yield(bad, nil) {
return
}
yield(good, nil)
}
it := newElemIter(seq, staging, false)
defer func() { _ = it.Close() }()
if !it.Next(t.Context()) {
t.Fatal("an unsafe name ended the walk; the item after it was never reached")
}
if got := it.current.item.MessageID; got != 8 {
t.Fatalf("Next yielded message %d, want the item after the unsafe one", got)
}
if it.Next(t.Context()) {
t.Error("iterator produced a third item")
}
if it.failure != nil {
t.Errorf("a skipped item must not fail the run, got: %v", it.failure)
}
if len(it.skipped) != 1 {
t.Fatalf("skipped = %d, want 1", len(it.skipped))
}
if !strings.Contains(it.skipped[0].Error(), "message 7") {
t.Errorf("the skip should name the message, got: %v", it.skipped[0])
}
}
// The skipped items still have to reach the caller: the run did work, but these
// messages were never attempted and nothing else would say so.
func TestDownloadReportsSkippedItems(t *testing.T) {
staging := t.TempDir()
bad := testItem(t, 7, "../../escape.conf", 10)
@@ -139,14 +177,14 @@ func TestElemIterRejectsUnsafeNames(t *testing.T) {
it := newElemIter(seq, staging, false)
defer func() { _ = it.Close() }()
if it.Next(t.Context()) {
t.Fatal("iterator accepted a name that escapes the staging directory")
for it.Next(t.Context()) {
}
if it.failure == nil {
t.Fatal("no failure recorded after rejecting an unsafe name")
err := errors.Join(append([]error{nil}, it.skipped...)...)
if err == nil {
t.Fatal("skipped items produced no error for the caller")
}
if !strings.Contains(it.failure.Error(), "message 7") {
t.Errorf("failure should name the message, got: %v", it.failure)
if !strings.Contains(err.Error(), "message 7") {
t.Errorf("error should name the skipped message, got: %v", err)
}
}
@@ -306,12 +344,6 @@ func TestElemIterNeverReportsErrToTheDownloader(t *testing.T) {
}
return newElemIter(seq, staging, false)
},
"unsafe name": func() *elemIter {
seq := func(yield func(tgsource.Item, error) bool) {
yield(testItem(t, 1, "../escape", 10), nil)
}
return newElemIter(seq, staging, false)
},
"cancelled": func() *elemIter {
seq := func(yield func(tgsource.Item, error) bool) {
yield(testItem(t, 2, "a.mp4", 10), nil)
+38 -21
View File
@@ -87,6 +87,11 @@ type elemIter struct {
// another goroutine without racing on the iterator itself.
stopped *atomic.Bool
// skipped collects items refused before any download was attempted. They do
// not stop the run: one message with a hostile filename must not be able to
// strand every message behind it, which is what ending iteration would mean.
skipped []error
// opened records every file handle. finish closes each one on the normal
// path, so this is not what keeps descriptors from leaking; it is the
// backstop for items that were opened but never reached finish, which is
@@ -106,29 +111,41 @@ func newElemIter(seq iter.Seq2[tgsource.Item, error], staging string, takeout bo
}
func (i *elemIter) Next(ctx context.Context) bool {
if i.failure != nil || i.stopped.Load() {
return false
}
if err := ctx.Err(); err != nil {
i.failure = err
return false
}
var item tgsource.Item
for {
if i.failure != nil || i.stopped.Load() {
return false
}
if err := ctx.Err(); err != nil {
i.failure = err
return false
}
item, err, ok := i.next()
if !ok {
return false
}
if err != nil {
i.failure = err
return false
}
var (
err error
ok bool
)
item, err, ok = i.next()
if !ok {
return false
}
if err != nil {
i.failure = err
return false
}
// A name that cannot be written is refused here rather than left to
// os.Create: the error names the message, and the run continues instead of
// failing on a path that could never have worked.
if err := naming.Safe(item.Name); err != nil {
i.failure = fmt.Errorf("message %d: %w", item.MessageID, err)
return false
// A name that cannot be written is refused here rather than left to
// os.Create, and refusing it skips the item rather than ending the walk.
// selectTodo already filters these out on the CLI path, so reaching this
// is either a second caller or a gap there; in both cases one unwritable
// name must not strand the rest of the chat behind it.
if err := naming.Safe(item.Name); err != nil {
if len(i.skipped) < maxRecordedErrors {
i.skipped = append(i.skipped, fmt.Errorf("message %d: %w", item.MessageID, err))
}
continue
}
break
}
if i.acquire != nil {
+19 -7
View File
@@ -61,6 +61,14 @@ func (r Result) Failed() []Outcome {
// makes the final message unreadable.
const maxRecordedErrors = 10
// download is the download step, indirected so a test can drive Run's
// composition without a live Telegram connection. What that buys is coverage of
// the three properties Run alone is responsible for — that the upload channel is
// closed only after every send, that the byte budget balances across a whole
// run, and that a tripped breaker still terminates — none of which the pieces
// can be tested for individually.
var download = Download
// Run downloads every item and uploads each one as it completes.
//
// This is the whole reason for the rewrite. run.sh could not see inside tdl, so
@@ -150,7 +158,7 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R
}()
}
dlOutcomes, stats, dlErr := Download(ctx, seq, DownloadOptions{
dlOutcomes, stats, dlErr := download(ctx, seq, DownloadOptions{
Pool: o.Pool,
Staging: o.Staging,
Threads: o.Threads,
@@ -178,17 +186,21 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R
trip, total := tripped, nErrs
mu.Unlock()
res := Result{Stats: stats, Outcomes: dlOutcomes}
// Both halves are reported. On the most common failure path — Ctrl-C — the
// download side returns context.Canceled while the upload workers drain
// whatever is still queued, fail every one of them against the cancelled
// context, and delete the staged file each time. Returning only the download
// error would leave the operator with "context canceled" and no sign that
// finished files had been discarded.
var upErr error
switch {
case dlErr != nil:
return res, dlErr
case trip:
return res, fmt.Errorf("stopped after %d consecutive upload failures (%d total): %w",
upErr = fmt.Errorf("stopped after %d consecutive upload failures (%d total): %w",
o.MaxFailures, total, joined)
case total > 0:
return res, fmt.Errorf("%d upload(s) failed: %w", total, joined)
upErr = fmt.Errorf("%d upload(s) failed: %w", total, joined)
}
return res, nil
return Result{Stats: stats, Outcomes: dlOutcomes}, errors.Join(dlErr, upErr)
}
// budget bounds how many bytes of downloaded-but-not-yet-uploaded data sit on
+306
View File
@@ -0,0 +1,306 @@
package pipeline
import (
"context"
"errors"
"fmt"
"iter"
"os"
"path/filepath"
"strings"
"sync/atomic"
"testing"
_ "github.com/rclone/rclone/backend/local"
"github.com/rclone/rclone/fs"
"github.com/iyear/tdl/core/tmedia"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// These exercise Run's composition rather than its parts. Everything underneath
// it is tested directly, but the properties Run alone owns — that the upload
// channel closes only after the last send, that reservations balance across a
// whole run, that a tripped breaker still terminates — only exist once the
// pieces are wired together, and a regression in any of them is silent.
// runFake substitutes the download step for the duration of a test.
func runFake(t *testing.T, f func(context.Context, iter.Seq2[tgsource.Item, error], DownloadOptions) ([]Outcome, Stats, error)) {
t.Helper()
prev := download
download = f
t.Cleanup(func() { download = prev })
}
// stageItems is a download step that writes each item's bytes into staging and
// hands it to the upload leg, mimicking what the real one does on success.
// Items named in fail never reach staging.
func stageItems(fail map[int]bool) func(context.Context, iter.Seq2[tgsource.Item, error], DownloadOptions) ([]Outcome, Stats, error) {
return func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
var (
outcomes []Outcome
stats Stats
)
for it, err := range seq {
if err != nil {
return outcomes, stats, err
}
if o.stop != nil && o.stop.Load() {
break
}
if o.acquire != nil {
if aerr := o.acquire(ctx, it.Size()); aerr != nil {
return outcomes, stats, aerr
}
}
stats.Started++
if fail[it.MessageID] {
stats.Failed++
outcomes = append(outcomes, Outcome{Item: it, Err: errors.New("download failed")})
o.onFailed(it)
continue
}
if werr := os.WriteFile(filepath.Join(o.Staging, it.Name),
make([]byte, it.Size()), 0o600); werr != nil {
return outcomes, stats, werr
}
stats.Done++
stats.BytesDone += it.Size()
outcomes = append(outcomes, Outcome{Item: it})
o.onReady(it)
}
return outcomes, stats, nil
}
}
func testItems(n int, size int64) []tgsource.Item {
out := make([]tgsource.Item, n)
for i := range out {
out[i] = tgsource.Item{
MessageID: i + 1,
Name: fmt.Sprintf("-100123_%d_file.bin", i+1),
Media: &tmedia.Media{Size: size},
}
}
return out
}
func seqOf(items []tgsource.Item) iter.Seq2[tgsource.Item, error] {
return func(yield func(tgsource.Item, error) bool) {
for _, it := range items {
if !yield(it, nil) {
return
}
}
}
}
// runOpts wires Run against local directories, so uploads are real rclone moves.
func runOpts(t *testing.T, budget int64) (Options, string, string) {
t.Helper()
staging, dstDir := t.TempDir(), t.TempDir()
dst, err := fs.NewFs(t.Context(), dstDir)
if err != nil {
t.Fatalf("open destination: %v", err)
}
return Options{
Dst: dst,
Staging: staging,
Uploads: 2,
Budget: budget,
Confirm: true,
}, staging, dstDir
}
func TestRunUploadsEveryDownloadedItem(t *testing.T) {
runFake(t, stageItems(nil))
o, staging, dstDir := runOpts(t, 0)
items := testItems(6, 512)
res, err := Run(t.Context(), seqOf(items), o)
if err != nil {
t.Fatalf("Run: %v", err)
}
if res.Stats.Done != len(items) {
t.Errorf("Done = %d, want %d", res.Stats.Done, len(items))
}
for _, it := range items {
info, serr := os.Stat(filepath.Join(dstDir, it.Name))
if serr != nil {
t.Errorf("%s not on the destination: %v", it.Name, serr)
continue
}
if info.Size() != it.Size() {
t.Errorf("%s is %d bytes, want %d", it.Name, info.Size(), it.Size())
}
}
assertEmpty(t, staging)
}
// The reason Err always reports nil: if Run closed the upload channel before
// the download step finished handing over items, this panics.
func TestRunClosesUploadsOnlyAfterTheLastSend(t *testing.T) {
var late atomic.Bool
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
// A worker still delivering after the step's own error is exactly what
// core does when it skips wg.Wait, so send one and then fail.
it := testItems(1, 128)[0]
if err := os.WriteFile(filepath.Join(o.Staging, it.Name), make([]byte, it.Size()), 0o600); err != nil {
return nil, Stats{}, err
}
o.acquire(ctx, it.Size())
o.onReady(it)
late.Store(true)
return []Outcome{{Item: it}}, Stats{Started: 1, Done: 1}, errors.New("download step failed")
})
o, staging, dstDir := runOpts(t, 0)
_, err := Run(t.Context(), seqOf(nil), o)
if err == nil {
t.Fatal("Run returned nil, want the download step's error")
}
if !late.Load() {
t.Fatal("the download step never ran")
}
// The handed-over item must still have been uploaded, not dropped.
if _, serr := os.Stat(filepath.Join(dstDir, "-100123_1_file.bin")); serr != nil {
t.Errorf("item handed over before the error was not uploaded: %v", serr)
}
assertEmpty(t, staging)
}
// A reservation that is not returned shrinks the cap for the rest of the run,
// and one returned twice panics. Neither is visible in the pieces individually.
//
// The budget here is exactly one file, so the run can only proceed if every
// reservation comes back: the second item cannot start until the first is
// released. A leak deadlocks and this test times out rather than passing
// quietly, which is the whole point of sizing it this way.
func TestRunBalancesTheBudgetAcrossFailures(t *testing.T) {
const size = 1024
runFake(t, stageItems(map[int]bool{2: true, 5: true}))
o, staging, dstDir := runOpts(t, size)
items := testItems(8, size)
res, err := Run(t.Context(), seqOf(items), o)
if err != nil {
t.Fatalf("Run: %v", err)
}
if got := len(res.Failed()); got != 2 {
t.Errorf("Failed() = %d, want 2", got)
}
if res.Stats.Done != 6 {
t.Errorf("Done = %d, want 6", res.Stats.Done)
}
// The six that downloaded are on the destination; the two that failed are
// not, and neither is holding space.
for _, it := range items {
_, serr := os.Stat(filepath.Join(dstDir, it.Name))
wantThere := it.MessageID != 2 && it.MessageID != 5
if wantThere && serr != nil {
t.Errorf("%s should be on the destination: %v", it.Name, serr)
}
if !wantThere && !os.IsNotExist(serr) {
t.Errorf("%s should not be on the destination", it.Name)
}
}
assertEmpty(t, staging)
}
func TestRunStopsDownloadingAfterConsecutiveUploadFailures(t *testing.T) {
const size = 256
items := testItems(40, size)
var dispatched atomic.Int64
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
var (
outcomes []Outcome
stats Stats
)
for it := range seqValues(seq) {
if o.stop.Load() {
break
}
o.acquire(ctx, it.Size())
dispatched.Add(1)
// Nothing is written to staging, so every upload fails to find it.
stats.Started++
stats.Done++
outcomes = append(outcomes, Outcome{Item: it})
o.onReady(it)
}
return outcomes, stats, nil
})
o, staging, _ := runOpts(t, 0)
o.MaxFailures = 3
_, err := Run(t.Context(), seqOf(items), o)
if err == nil {
t.Fatal("Run returned nil, want the breaker's error")
}
if !strings.Contains(err.Error(), "consecutive upload failures") {
t.Errorf("error does not mention the breaker: %v", err)
}
// The point of stopping the iterator rather than cancelling uploads: the
// download side must not have walked the whole chat.
if got := dispatched.Load(); got == int64(len(items)) {
t.Errorf("all %d items were dispatched; the breaker did not stop downloads", got)
}
assertEmpty(t, staging)
}
func seqValues(seq iter.Seq2[tgsource.Item, error]) iter.Seq[tgsource.Item] {
return func(yield func(tgsource.Item) bool) {
for it, err := range seq {
if err != nil {
return
}
if !yield(it) {
return
}
}
}
}
// A failed upload must take the staged file with it, or the cap is over-committed
// by that much for the rest of the run.
func TestRunClearsStagingWhenAnUploadFails(t *testing.T) {
const size = 512
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
it := testItems(1, size)[0]
// Staged at the wrong size, so the confirm step rejects it.
if err := os.WriteFile(filepath.Join(o.Staging, it.Name), make([]byte, size/2), 0o600); err != nil {
return nil, Stats{}, err
}
o.acquire(ctx, it.Size())
o.onReady(it)
return []Outcome{{Item: it}}, Stats{Started: 1, Done: 1}, nil
})
o, staging, dstDir := runOpts(t, 0)
_, err := Run(t.Context(), seqOf(nil), o)
if err == nil {
t.Fatal("Run returned nil, want the confirm failure")
}
assertEmpty(t, staging)
// And the short object must not be left under the name verify matches.
if _, serr := os.Stat(filepath.Join(dstDir, "-100123_1_file.bin")); !os.IsNotExist(serr) {
t.Errorf("short object left on the destination: %v", serr)
}
}
func assertEmpty(t *testing.T, dir string) {
t.Helper()
entries, err := os.ReadDir(dir)
if err != nil {
t.Fatalf("read staging: %v", err)
}
if len(entries) != 0 {
names := make([]string, len(entries))
for i, e := range entries {
names[i] = e.Name()
}
t.Errorf("staging is not empty: %v", names)
}
}
+53 -14
View File
@@ -2,6 +2,7 @@ package pipeline
import (
"context"
"errors"
"fmt"
"time"
@@ -24,16 +25,22 @@ type uploader struct {
// MoveFile removes the local copy as part of the move, so a successful return
// means the file is on the remote and off local disk.
//
// A short object is deleted rather than left in place, and that is the part that
// matters. Verification matches on name and non-zero size, so a truncated object
// under the right name would be counted archived by this run and by every run
// after it — permanently, with the local copy already gone. Removing it turns a
// silent corruption into an absent file the next run fetches again. The shell
// pipeline had this hole too: it noticed a bad upload only at the next verify,
// by which point the evidence was the same.
// Both failure paths end at dropShort, and that is the part that matters. An
// object of the wrong size sitting under the right name is worse than no object
// at all: verification matches on name and non-zero size, so it would be counted
// archived by this run and by every run after it — permanently, once the local
// copy is gone. Removing it turns a silent corruption into an absent file the
// next run fetches again.
func (u *uploader) upload(ctx context.Context, it tgsource.Item) error {
if err := operations.MoveFile(ctx, u.dst, u.local, it.Name, it.Name); err != nil {
return fmt.Errorf("move %q to %s: %w", it.Name, u.dst.String(), err)
err = fmt.Errorf("move %q to %s: %w", it.Name, u.dst.String(), err)
// A failed move can still leave a partial object under the final name.
// rclone only writes to a temporary name when the backend advertises
// PartialUploads (copy.go:93), and it only cleans up after itself when
// it did (copy.go:348-350) — pikpak, the remote this was built against,
// advertises neither, so a died-halfway transfer stays exactly where a
// complete one would be.
return errors.Join(err, u.dropShort(ctx, it))
}
if !u.confirm {
return nil
@@ -45,12 +52,7 @@ func (u *uploader) upload(ctx context.Context, it tgsource.Item) error {
}
if got := obj.Size(); got != it.Size() {
err := fmt.Errorf("confirm %q: remote has %d bytes, expected %d", it.Name, got, it.Size())
// Deleted on a fresh context: the run may already be shutting down, and
// leaving a plausible-looking short object behind is worse than the
// error that got us here.
delCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if derr := operations.DeleteFile(delCtx, obj); derr != nil {
if derr := remove(ctx, obj); derr != nil {
return fmt.Errorf("%w (and it could not be removed: %v — delete it by hand "+
"or verify will count it archived)", err, derr)
}
@@ -58,3 +60,40 @@ func (u *uploader) upload(ctx context.Context, it tgsource.Item) error {
}
return nil
}
// dropShort removes an object left under it.Name at the wrong size.
//
// An object of the *right* size is deliberately left alone. pikpak commits an
// upload as a server-side async task, so a transfer rclone gave up on can still
// land correctly afterwards; deleting it on the strength of the error alone
// would throw away a good file and force it to be fetched again.
func (u *uploader) dropShort(ctx context.Context, it tgsource.Item) error {
// A fresh context: the run may already be shutting down, which is one of
// the ways the move failed in the first place.
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
obj, err := u.dst.NewObject(ctx, it.Name)
if errors.Is(err, fs.ErrorObjectNotFound) {
return nil // nothing was left behind
}
if err != nil {
return fmt.Errorf("check for a leftover %q: %w", it.Name, err)
}
if obj.Size() == it.Size() {
return nil
}
if derr := remove(ctx, obj); derr != nil {
return fmt.Errorf("a %d-byte fragment of %q (expected %d) is on the remote and "+
"could not be removed: %w — delete it by hand or verify will count it archived",
obj.Size(), it.Name, it.Size(), derr)
}
return nil
}
// remove deletes an object on a context that outlives the run's cancellation.
func remove(ctx context.Context, obj fs.Object) error {
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
return operations.DeleteFile(ctx, obj)
}
+68
View File
@@ -0,0 +1,68 @@
package pipeline
import (
"os"
"path/filepath"
"testing"
_ "github.com/rclone/rclone/backend/local"
"github.com/rclone/rclone/fs"
"github.com/iyear/tdl/core/tmedia"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// dropShort is what stands between a died-halfway upload and a permanently
// wrong archive. rclone writes straight to the final name on any backend that
// does not advertise PartialUploads (copy.go:93) and cleans up only when it did
// not (copy.go:348-350), so on such a remote — pikpak, here — a fragment is left
// under exactly the name verification matches on.
func TestDropShortRemovesAFragmentButKeepsACompleteFile(t *testing.T) {
const name = "-100123_4242_clip.mp4"
const want = 4096
cases := map[string]struct {
staged int // bytes already at the destination; -1 means absent
wantThere bool
}{
// A fragment must go: left alone, verify counts it archived by name and
// non-zero size, for this run and every run after it.
"fragment is removed": {staged: 400, wantThere: false},
// A complete file must stay. pikpak commits uploads as a server-side
// async task, so a transfer rclone gave up on can still land correctly;
// deleting on the strength of the error alone throws away a good file.
"complete file is kept": {staged: want, wantThere: true},
"nothing to clean up": {staged: -1, wantThere: false},
}
for label, tc := range cases {
t.Run(label, func(t *testing.T) {
dstDir := t.TempDir()
path := filepath.Join(dstDir, name)
if tc.staged >= 0 {
if err := os.WriteFile(path, make([]byte, tc.staged), 0o600); err != nil {
t.Fatalf("stage destination file: %v", err)
}
}
dst, err := fs.NewFs(t.Context(), dstDir)
if err != nil {
t.Fatalf("open destination: %v", err)
}
u := &uploader{dst: dst}
it := tgsource.Item{MessageID: 4242, Name: name, Media: &tmedia.Media{Size: want}}
if err := u.dropShort(t.Context(), it); err != nil {
t.Fatalf("dropShort: %v", err)
}
_, serr := os.Stat(path)
switch {
case tc.wantThere && serr != nil:
t.Errorf("a complete file was deleted: %v", serr)
case !tc.wantThere && serr == nil:
t.Error("a short object was left under the name verify matches")
}
})
}
}