diff --git a/internal/pipeline/download.go b/internal/pipeline/download.go index 5e4a0ba..7d9bf8a 100644 --- a/internal/pipeline/download.go +++ b/internal/pipeline/download.go @@ -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 diff --git a/internal/pipeline/download_test.go b/internal/pipeline/download_test.go index 9952a5b..691fbdb 100644 --- a/internal/pipeline/download_test.go +++ b/internal/pipeline/download_test.go @@ -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) diff --git a/internal/pipeline/elem.go b/internal/pipeline/elem.go index 703abe8..4738202 100644 --- a/internal/pipeline/elem.go +++ b/internal/pipeline/elem.go @@ -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 { diff --git a/internal/pipeline/pipeline.go b/internal/pipeline/pipeline.go index af9b391..63a2487 100644 --- a/internal/pipeline/pipeline.go +++ b/internal/pipeline/pipeline.go @@ -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 diff --git a/internal/pipeline/run_test.go b/internal/pipeline/run_test.go new file mode 100644 index 0000000..e0841e3 --- /dev/null +++ b/internal/pipeline/run_test.go @@ -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) + } +} diff --git a/internal/pipeline/upload.go b/internal/pipeline/upload.go index f38f108..47670eb 100644 --- a/internal/pipeline/upload.go +++ b/internal/pipeline/upload.go @@ -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) +} diff --git a/internal/pipeline/upload_test.go b/internal/pipeline/upload_test.go new file mode 100644 index 0000000..c61f5ad --- /dev/null +++ b/internal/pipeline/upload_test.go @@ -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") + } + }) + } +}