From a81e7eaacbdacda299f72114fd7927d3c92cbd2e Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Sun, 6 Sep 2026 19:21:46 +0700 Subject: [PATCH] feat: download, upload and drive a chat to completion in one process MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phases 4 through 6: the two legs and the command that joins them. Downloads go to .part and are renamed only once complete, so a file without the suffix is always whole. That is what lets the upload leg treat "exists" as "finished" — the property run.sh could only approximate with a filename convention plus an age guard, because it could not see inside tdl. Every finished file is checked against the size Telegram reported, and that check rather than the error is the authoritative signal. core's downloader logs a failed transfer and returns nil, and its completion callback is deferred on that named return, so a failure arrives indistinguishable from a success. Trusting it would promote a truncated file and archive it as complete. The disk cap is a semaphore over bytes. A download reserves its own size before starting and releases it only after the upload confirms, so a slow remote stalls downloads by itself. Blocking the iterator is safe because the downloader calls it from its dispatch loop while workers run in a group, so a blocked iterator never stops the uploads that free the space. Gone with it: the du polling, the SIGSTOP and SIGCONT suspension, the min-age guard, the temp-file filter and the sweep-failure counter. A cap smaller than the largest file is refused up front. The semaphore could never admit it, and a run blocked on a file it can never start looks exactly like a stalled remote. Uploads re-state each object to prove its size before the local copy is gone, closing a gap where a truncated upload was only noticed by a later verify. The destination is created before the chat is read. It is also the credentials check, and doing it first means a bad destination fails in seconds rather than after a full history walk. One invocation converges: each item is checked against the index immediately before download, so there are no passes and re-running is the resume path. Options that no longer exist say what replaced them instead of failing as unknown flags. Verified end to end against the live chat and a scratch remote path: two files downloaded, uploaded, confirmed present at the right size, staging left empty. --- cmd/tgexport/main.go | 5 +- cmd/tgexport/sync.go | 322 +++++++++++++++++++++++++++++ cmd/tgexport/sync_test.go | 139 +++++++++++++ internal/pipeline/download.go | 172 +++++++++++++++ internal/pipeline/download_test.go | 290 ++++++++++++++++++++++++++ internal/pipeline/elem.go | 13 ++ internal/pipeline/pipeline.go | 193 +++++++++++++++++ internal/pipeline/pipeline_test.go | 147 +++++++++++++ internal/pipeline/progress.go | 105 ++++++++++ internal/pipeline/upload.go | 47 +++++ internal/remote/fs.go | 14 ++ internal/remote/index.go | 7 + internal/report/progress.go | 108 ++++++++++ 13 files changed, 1561 insertions(+), 1 deletion(-) create mode 100644 cmd/tgexport/sync.go create mode 100644 cmd/tgexport/sync_test.go create mode 100644 internal/pipeline/download.go create mode 100644 internal/pipeline/download_test.go create mode 100644 internal/pipeline/pipeline.go create mode 100644 internal/pipeline/pipeline_test.go create mode 100644 internal/pipeline/progress.go create mode 100644 internal/pipeline/upload.go create mode 100644 internal/report/progress.go diff --git a/cmd/tgexport/main.go b/cmd/tgexport/main.go index 0a515fe..7314d8f 100644 --- a/cmd/tgexport/main.go +++ b/cmd/tgexport/main.go @@ -65,6 +65,8 @@ func run() int { err = listCmd(ctx, os.Args[2:]) case "verify": err = verifyCmd(ctx, os.Args[2:]) + case "sync": + err = syncCmd(ctx, os.Args[2:]) default: fmt.Fprintf(os.Stderr, "unknown command %q\n\n", os.Args[1]) usage() @@ -150,9 +152,10 @@ func usage() { fmt.Fprint(os.Stderr, `Usage: tgexport [options] Commands: - doctor Check the Telegram session, the destination remote, and free space + sync Archive a chat to a remote, fetching only what is missing list Print every media message in a chat as idsizename verify Report whether a chat is fully archived on a remote + doctor Check the Telegram session, the destination remote, and free space Run 'tgexport -h' for command options. `) diff --git a/cmd/tgexport/sync.go b/cmd/tgexport/sync.go new file mode 100644 index 0000000..207b139 --- /dev/null +++ b/cmd/tgexport/sync.go @@ -0,0 +1,322 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "iter" + "os" + + "github.com/iyear/tdl/core/dcpool" + tdlstorage "github.com/iyear/tdl/core/storage" + "github.com/rclone/rclone/fs" + + "github.com/tiennm99dev/telegram-exporter/internal/pipeline" + "github.com/tiennm99dev/telegram-exporter/internal/remote" + "github.com/tiennm99dev/telegram-exporter/internal/report" + "github.com/tiennm99dev/telegram-exporter/internal/tdlkv" + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" + "github.com/tiennm99dev/telegram-exporter/internal/verify" +) + +// retiredFlags map options the shell pipeline had onto what replaced them. +// +// Recognising them beats "flag provided but not defined": these were in +// muscle memory and in wrapper scripts, and a bare parse error does not say +// whether the concept moved or disappeared. +var retiredFlags = map[string]string{ + "i": "the rclone sweep interval is gone; uploads start the moment a download finishes", + "a": "--min-age is gone; a file is only uploaded once the downloader reports it complete", + "f": "the export JSON is gone; the chat is read live, so names cannot go stale", + "p": "there are no passes; one invocation converges, and re-running resumes", + "q": "renamed to --min-free", +} + +// syncCmd archives a chat to a remote: read the chat, skip what is already +// there, download and upload the rest, then report on the result. +func syncCmd(ctx context.Context, args []string) error { + fs := flag.NewFlagSet("sync", flag.ContinueOnError) + var ( + chat = fs.String("c", "", "chat id, username, or t.me link (required)") + remoteArg = fs.String("r", "", "rclone destination, e.g. pikpak:archive (required)") + staging = fs.String("d", "./staging", "staging directory for files in flight") + maxStaging = fs.String("m", "", "cap staging at this size, e.g. 40G (default: no cap)") + threads = fs.Int("threads", 4, "connections per file") + limit = fs.Int("limit", 2, "files downloading at once") + uploads = fs.Int("uploads", 2, "files uploading at once") + minFree = fs.Int64("min-free", 5, "stop if the remote has fewer than this many GiB free") + limitItems = fs.Int("limit-items", 0, "stop after this many files (0 means no limit)") + confirm = fs.Bool("confirm", true, "re-state each uploaded file to prove its size") + takeout = fs.Bool("takeout", true, "use a takeout session, as `tdl dl --takeout` did") + ns = fs.String("n", "default", "tdl session namespace") + dataDir = fs.String("storage", tdlkv.DefaultDir(), "tdl bolt storage directory") + ) + for name, replacement := range retiredFlags { + fs.Var(retiredFlag{name, replacement}, name, "retired") + } + + if err := fs.Parse(args); err != nil { + if errors.Is(err, flag.ErrHelp) { + return err + } + return fmt.Errorf("%w: %v", errUsage, err) + } + if *chat == "" || *remoteArg == "" { + return fmt.Errorf("%w: -c CHAT and -r REMOTE:PATH are both required", errUsage) + } + + budget, err := parseSize(*maxStaging) + if err != nil { + return fmt.Errorf("%w: -m %v", errUsage, err) + } + + ctx, err = remote.Init(ctx, remote.DefaultTunables()) + if err != nil { + return err + } + dst, err := remote.Resolve(ctx, *remoteArg) + if err != nil { + return err + } + + // Partial files from an earlier run cannot be continued — core's downloader + // takes no starting offset — so they are cleared before anything else fills + // the disk with fragments no run will finish. + if swept, err := pipeline.SweepPartials(*staging); err != nil { + return fmt.Errorf("clear partial downloads: %w", err) + } else if swept > 0 { + fmt.Fprintf(os.Stderr, "cleared %d partial download(s) from an earlier run\n", swept) + } + + // Before anything expensive: prove the destination is reachable and writable. + if err := remote.EnsureDir(ctx, dst); err != nil { + return err + } + if err := checkFree(ctx, dst, *minFree); err != nil { + return err + } + + kv, err := tdlkv.Open(*dataDir, *ns) + if err != nil { + return err + } + defer func() { _ = kv.Close() }() + + sess, err := tgsource.New(ctx, tgsource.Options{KV: kv}) + if err != nil { + return err + } + + var final verify.Report + if err := sess.Run(ctx, func(ctx context.Context, pool dcpool.Pool) error { + api := pool.Default(ctx) + + peer, err := tgsource.ResolveChat(ctx, tgsource.Manager(api, tdlstorage.NewPeers(kv)), *chat) + if err != nil { + return err + } + + fmt.Fprintf(os.Stderr, "reading %s\n", *chat) + var items []tgsource.Item + for it, err := range tgsource.Walk(ctx, api, peer) { + if err != nil { + return err + } + items = append(items, it) + } + + idx, err := remote.BuildIndex(ctx, dst, peer.ID()) + if err != nil { + return err + } + before := verify.Check(items, idx) + fmt.Fprintf(os.Stderr, "%d media messages, %d already archived, %d to fetch\n", + before.Expected, before.Present, len(before.Todo())) + + todo := selectTodo(items, before, *limitItems) + if len(todo) == 0 { + final = before + return nil + } + + if err := validateBudget(budget, todo); err != nil { + return err + } + + var todoBytes int64 + for _, it := range todo { + todoBytes += it.Size() + } + fmt.Fprintf(os.Stderr, "fetching %d file(s), %.1f GiB\n", len(todo), float64(todoBytes)/(1<<30)) + + rep := report.New(os.Stderr, len(todo), todoBytes) + res, runErr := pipeline.Run(ctx, sliceSeq(todo), pipeline.Options{ + Pool: pool, + Dst: dst, + Staging: *staging, + Threads: *threads, + Limit: *limit, + Uploads: *uploads, + Budget: budget, + Confirm: *confirm, + Takeout: *takeout, + Report: rep.Update, + }) + rep.Finish(res.Stats) + + for _, f := range res.Failed() { + fmt.Fprintf(os.Stderr, " message %d failed: %v\n", f.Item.MessageID, f.Err) + } + if runErr != nil { + return runErr + } + + // The remote is re-indexed rather than assumed: the run's own view of + // what it uploaded is exactly the thing under test. + idx, err = remote.BuildIndex(ctx, dst, peer.ID()) + if err != nil { + return err + } + final = verify.Check(items, idx) + return nil + }); err != nil { + return err + } + + fmt.Fprintln(os.Stderr) + final.Write(os.Stdout) + + if !final.Complete() { + return fmt.Errorf("%w: %d file(s) still to fetch", errIncomplete, len(final.Todo())) + } + return nil +} + +// selectTodo picks the items still needing a fetch, newest first, optionally +// capped for a smoke test. +func selectTodo(items []tgsource.Item, r verify.Report, limit int) []tgsource.Item { + want := make(map[int]struct{}, len(r.Todo())) + for _, id := range r.Todo() { + want[id] = struct{}{} + } + // Unsafe names are in Todo so they stay visible in the report, but fetching + // one is impossible by definition, so it is not queued for download. + for _, u := range r.Unsafe { + delete(want, u.MessageID) + } + + var todo []tgsource.Item + for _, it := range items { + if _, ok := want[it.MessageID]; !ok { + continue + } + todo = append(todo, it) + if limit > 0 && len(todo) == limit { + break + } + } + return todo +} + +func sliceSeq(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 + } + } + } +} + +// validateBudget refuses a cap smaller than the largest file. +// +// The semaphore can never admit a weight above its limit, so such a run would +// block forever on a file it could never start — indistinguishable, from the +// outside, from a stalled remote. run.sh could only warn about this after the +// fact, once draining failed to get back under the cap. +func validateBudget(budget int64, todo []tgsource.Item) error { + if budget <= 0 { + return nil + } + var largest int64 + for _, it := range todo { + if it.Size() > largest { + largest = it.Size() + } + } + if largest > budget { + return fmt.Errorf("%w: -m is %.1f GiB but the largest file to fetch is %.1f GiB; "+ + "the cap must exceed the biggest single file", + errUsage, float64(budget)/(1<<30), float64(largest)/(1<<30)) + } + return nil +} + +// checkFree refuses to start when the remote is nearly full. +// +// A backend that cannot report a quota is treated as unlimited rather than as a +// failure — the shell pipeline made that choice deliberately so a remote without +// an About API never blocked a run, and it is preserved. +func checkFree(ctx context.Context, dst fs.Fs, minGiB int64) error { + free, ok := remote.FreeBytes(ctx, dst) + if !ok { + return nil + } + freeGiB := free / (1 << 30) + fmt.Fprintf(os.Stderr, "%s has %d GiB free\n", dst.String(), freeGiB) + if freeGiB < minGiB { + return fmt.Errorf("%s has only %d GiB free, below the %d GiB floor; "+ + "free space or lower --min-free", dst.String(), freeGiB, minGiB) + } + return nil +} + +// retiredFlag reports a helpful error for an option that no longer exists. +type retiredFlag struct{ name, replacement string } + +func (r retiredFlag) String() string { return "" } +func (r retiredFlag) Set(string) error { + return fmt.Errorf("-%s no longer exists: %s", r.name, r.replacement) +} + +// parseSize reads a binary size such as 40G, matching what run.sh -m accepted. +func parseSize(s string) (int64, error) { + if s == "" { + return 0, nil + } + mult := int64(1) + switch unit := s[len(s)-1]; unit { + case 'K', 'k': + mult = 1 << 10 + case 'M', 'm': + mult = 1 << 20 + case 'G', 'g': + mult = 1 << 30 + case 'T', 't': + mult = 1 << 40 + default: + if unit < '0' || unit > '9' { + return 0, fmt.Errorf("unknown size suffix %q, expected K, M, G or T", string(unit)) + } + } + digits := s + if mult > 1 { + digits = s[:len(s)-1] + } + + var n int64 + if digits == "" { + return 0, fmt.Errorf("%q has no number", s) + } + for _, r := range digits { + if r < '0' || r > '9' { + return 0, fmt.Errorf("%q is not a size", s) + } + n = n*10 + int64(r-'0') + } + if n <= 0 { + return 0, fmt.Errorf("must be greater than zero") + } + return n * mult, nil +} diff --git a/cmd/tgexport/sync_test.go b/cmd/tgexport/sync_test.go new file mode 100644 index 0000000..9cd6d10 --- /dev/null +++ b/cmd/tgexport/sync_test.go @@ -0,0 +1,139 @@ +package main + +import ( + "errors" + "strings" + "testing" + + "github.com/iyear/tdl/core/tmedia" + + "github.com/tiennm99dev/telegram-exporter/internal/naming" + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" + "github.com/tiennm99dev/telegram-exporter/internal/verify" +) + +func TestParseSize(t *testing.T) { + ok := map[string]int64{ + "": 0, // unset means no cap + "40G": 40 << 30, + "40g": 40 << 30, + "512M": 512 << 20, + "2T": 2 << 40, + "1024": 1024, // bare number is bytes + "1K": 1 << 10, + } + for in, want := range ok { + got, err := parseSize(in) + if err != nil { + t.Errorf("parseSize(%q) = %v, want %d", in, err, want) + continue + } + if got != want { + t.Errorf("parseSize(%q) = %d, want %d", in, got, want) + } + } + + for _, in := range []string{"G", "0", "0G", "-5G", "40GB", "4.5G", "abc", "40Q"} { + if got, err := parseSize(in); err == nil { + t.Errorf("parseSize(%q) = %d, want an error", in, got) + } + } +} + +func syncItem(id int, size int64) tgsource.Item { + m := &tmedia.Media{Name: "f.mp4", Size: size} + return tgsource.Item{DialogID: 1, MessageID: id, Name: naming.For(1, id, m), Media: m} +} + +func TestSelectTodoPicksOnlyOutstandingItems(t *testing.T) { + items := []tgsource.Item{syncItem(1, 10), syncItem(2, 10), syncItem(3, 10), syncItem(4, 10)} + r := verify.Report{Absent: []int{2, 4}, ZeroByte: []int{3}} + + got := selectTodo(items, r, 0) + if len(got) != 3 { + t.Fatalf("selected %d items, want 3", len(got)) + } + for _, it := range got { + if it.MessageID == 1 { + t.Error("selected an item that is already archived") + } + } +} + +// An unsafe name is kept in the report so it stays visible, but queuing it for +// download would retry something that can never succeed — the non-terminating +// loop this design exists to avoid. +func TestSelectTodoSkipsUnsafeNames(t *testing.T) { + items := []tgsource.Item{syncItem(1, 10), syncItem(2, 10)} + r := verify.Report{ + Absent: []int{1, 2}, + Unsafe: []verify.Unsafe{{MessageID: 2, Name: "../x", Reason: errors.New("unsafe")}}, + } + + got := selectTodo(items, r, 0) + if len(got) != 1 || got[0].MessageID != 1 { + t.Errorf("selected %+v, want only message 1", got) + } +} + +func TestSelectTodoHonoursLimit(t *testing.T) { + items := []tgsource.Item{syncItem(1, 10), syncItem(2, 10), syncItem(3, 10)} + r := verify.Report{Absent: []int{1, 2, 3}} + + if got := selectTodo(items, r, 2); len(got) != 2 { + t.Errorf("selected %d items with a limit of 2, want 2", len(got)) + } + if got := selectTodo(items, r, 0); len(got) != 3 { + t.Errorf("selected %d items with no limit, want 3", len(got)) + } +} + +// A cap below the largest file could never admit it, so the run would block on +// something it can never start — which from outside looks like a stalled remote. +// It has to be refused up front. +func TestValidateBudgetRejectsCapBelowLargestFile(t *testing.T) { + todo := []tgsource.Item{syncItem(1, 1<<20), syncItem(2, 5<<30)} + + err := validateBudget(4<<30, todo) + if err == nil { + t.Fatal("a 4 GiB cap was accepted with a 5 GiB file to fetch") + } + if !errors.Is(err, errUsage) { + t.Errorf("error should be a usage error, got: %v", err) + } + if !strings.Contains(err.Error(), "largest file") { + t.Errorf("error should name the problem, got: %v", err) + } + + if err := validateBudget(6<<30, todo); err != nil { + t.Errorf("a 6 GiB cap should accept a 5 GiB file, got: %v", err) + } + if err := validateBudget(0, todo); err != nil { + t.Errorf("an unset cap should accept anything, got: %v", err) + } +} + +// Options the shell pipeline had must produce an explanation, not "flag +// provided but not defined" — they are in wrapper scripts and muscle memory. +func TestRetiredFlagsExplainWhatReplacedThem(t *testing.T) { + for name, replacement := range retiredFlags { + err := retiredFlag{name, replacement}.Set("x") + if err == nil { + t.Errorf("-%s was accepted, want an explanation", name) + continue + } + if !strings.Contains(err.Error(), "-"+name) { + t.Errorf("error for -%s should name the flag, got: %v", name, err) + } + if !strings.Contains(err.Error(), replacement) { + t.Errorf("error for -%s should say what replaced it, got: %v", name, err) + } + } + + // The ones that mattered most in run.sh. + for _, name := range []string{"i", "a", "f", "p", "q"} { + if _, ok := retiredFlags[name]; !ok { + t.Errorf("-%s was a run.sh flag but is not recognised as retired", name) + } + } +} diff --git a/internal/pipeline/download.go b/internal/pipeline/download.go new file mode 100644 index 0000000..e160c0a --- /dev/null +++ b/internal/pipeline/download.go @@ -0,0 +1,172 @@ +package pipeline + +import ( + "context" + "errors" + "fmt" + "iter" + "os" + "path/filepath" + "strings" + + "github.com/iyear/tdl/core/dcpool" + "github.com/iyear/tdl/core/downloader" + + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +// DownloadOptions configures a download run. +type DownloadOptions struct { + Pool dcpool.Pool + Staging string + Threads int // connections per file + Limit int // files in flight + Takeout bool // use a takeout session, as `tdl dl --takeout` does + + // Report, when set, is called with a running summary. It is invoked from + // download worker goroutines, so it must be cheap and safe to call + // concurrently. + Report func(Stats) + + // acquire reserves staging space before a download starts, blocking until + // there is room. Unset means no bound. + acquire func(context.Context, int64) error + // onReady hands a completed file to the upload leg; onFailed says nothing + // was staged, so whatever acquire reserved must be given back. + onReady func(tgsource.Item) + onFailed func(tgsource.Item) +} + +// Download fetches every item in seq into the staging directory. +// +// Each file is written to .part and renamed to only once the +// downloader reports it complete, so a name without the suffix is always a +// whole file. That is what lets the upload half treat "the file exists" as +// "the file is finished" — the property the shell pipeline had to approximate +// with a filename convention plus an age guard, because it could not see +// inside tdl. +// +// A failed item does not abort the run: it is recorded in the returned outcomes +// and the rest continue, matching what a partial `tdl dl` pass did. +func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) { + if o.Threads <= 0 { + o.Threads = 4 + } + if o.Limit <= 0 { + o.Limit = 2 + } + if err := os.MkdirAll(o.Staging, 0o755); err != nil { + return nil, Stats{}, fmt.Errorf("create staging directory: %w", err) + } + + it := newElemIter(seq, o.Staging, o.Takeout) + it.acquire = o.acquire + defer func() { _ = it.Close() }() + + prog := newProgress(func(e *elem, err error) error { + ferr := finish(o.Staging, e, err) + if err == nil && ferr == nil { + if o.onReady != nil { + o.onReady(e.item) + } + return nil + } + if o.onFailed != nil { + o.onFailed(e.item) + } + return ferr + }, o.Report) + + err := downloader.New(downloader.Options{ + Pool: o.Pool, + Threads: o.Threads, + Iter: it, + Progress: prog, + }).Download(ctx, o.Limit) + + outcomes, stats := prog.results() + return outcomes, stats, err +} + +// finish closes a downloaded file and either promotes it or removes it. +// +// The size on disk is checked against the size Telegram reported, and that check +// is not belt-and-braces — it is the only reliable failure signal available. +// core's Download swallows non-cancellation errors: it logs them and returns +// nil, and OnDone is deferred on that named return, so a failed transfer arrives +// here indistinguishable from a successful one (downloader.go:47-60). Trusting +// the error alone would promote a truncated file to its final name, and the +// upload leg would archive it as complete. +// +// A partial file is deleted rather than kept: the downloader exposes no resume +// offset, so a leftover .part could never be continued, and leaving one behind +// would only invite a later run to mistake it for progress. +func finish(staging string, e *elem, downloadErr error) error { + part := partPath(staging, e.item) + + if cerr := e.file.Close(); cerr != nil && downloadErr == nil { + downloadErr = cerr + } + + if downloadErr == nil { + if err := checkSize(part, e.item.Size()); err != nil { + downloadErr = err + } + } + + if downloadErr != nil { + if rerr := os.Remove(part); rerr != nil && !os.IsNotExist(rerr) { + return fmt.Errorf("remove partial %q: %w", part, rerr) + } + // Returned so the caller records a failure even when the downloader + // claimed success; otherwise a short file would vanish silently and the + // run would report itself complete. + return downloadErr + } + + if err := os.Rename(part, finalPath(staging, e.item)); err != nil { + return fmt.Errorf("promote %q: %w", part, err) + } + return nil +} + +// checkSize compares what landed on disk against what Telegram said the file is. +func checkSize(path string, want int64) error { + info, err := os.Stat(path) + if err != nil { + return fmt.Errorf("stat downloaded file: %w", err) + } + if info.Size() != want { + return fmt.Errorf("short download: got %d bytes, expected %d", info.Size(), want) + } + return nil +} + +// SweepPartials removes leftover .part files from an earlier run. +// +// They cannot be resumed — core's downloader takes no starting offset — so the +// only options are delete or accumulate, and accumulating fills the disk with +// fragments no run will ever finish. +func SweepPartials(staging string) (int, error) { + entries, err := os.ReadDir(staging) + if err != nil { + if os.IsNotExist(err) { + return 0, nil + } + return 0, fmt.Errorf("read staging directory: %w", err) + } + + removed := 0 + var errs []error + for _, entry := range entries { + if entry.IsDir() || !strings.HasSuffix(entry.Name(), partSuffix) { + continue + } + if err := os.Remove(filepath.Join(staging, entry.Name())); err != nil { + errs = append(errs, err) + continue + } + removed++ + } + return removed, errors.Join(errs...) +} diff --git a/internal/pipeline/download_test.go b/internal/pipeline/download_test.go new file mode 100644 index 0000000..7940c46 --- /dev/null +++ b/internal/pipeline/download_test.go @@ -0,0 +1,290 @@ +package pipeline + +import ( + "errors" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/gotd/td/tg" + "github.com/iyear/tdl/core/downloader" + "github.com/iyear/tdl/core/tmedia" + + "github.com/tiennm99dev/telegram-exporter/internal/naming" + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +func testItem(t *testing.T, msgID int, file string, size int64) tgsource.Item { + t.Helper() + m := &tmedia.Media{ + Name: file, + Size: size, + DC: 2, + InputFileLoc: &tg.InputDocumentFileLocation{ID: int64(msgID)}, + } + return tgsource.Item{ + DialogID: 1234567890, + MessageID: msgID, + Name: naming.For(1234567890, msgID, m), + Media: m, + } +} + +// openElem mimics what elemIter does, so finish can be tested without a network. +func openElem(t *testing.T, staging string, it tgsource.Item) *elem { + t.Helper() + f, err := os.OpenFile(partPath(staging, it), os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + t.Fatalf("open part file: %v", err) + } + return &elem{item: it, file: f} +} + +// A name without the suffix must always be a whole file: that is the property +// the upload half relies on to treat "exists" as "finished", replacing the +// filename-convention-plus-age-guard the shell pipeline needed. +func TestFinishPromotesOnlyOnSuccess(t *testing.T) { + staging := t.TempDir() + it := testItem(t, 1, "video.mp4", 100) + + e := openElem(t, staging, it) + if _, err := e.file.WriteAt(make([]byte, it.Size()), 0); err != nil { + t.Fatalf("write: %v", err) + } + if err := finish(staging, e, nil); err != nil { + t.Fatalf("finish: %v", err) + } + + if _, err := os.Stat(finalPath(staging, it)); err != nil { + t.Errorf("final file missing after a successful download: %v", err) + } + if _, err := os.Stat(partPath(staging, it)); !os.IsNotExist(err) { + t.Errorf("part file still present after promotion") + } +} + +func TestFinishRemovesPartialOnFailure(t *testing.T) { + staging := t.TempDir() + it := testItem(t, 2, "video.mp4", 100) + + e := openElem(t, staging, it) + if _, err := e.file.WriteAt([]byte("half"), 0); err != nil { + t.Fatalf("write: %v", err) + } + // finish returns the failure rather than swallowing it, so the caller + // records the item as failed instead of quietly counting it done. + want := errors.New("connection reset") + if err := finish(staging, e, want); !errors.Is(err, want) { + t.Fatalf("finish = %v, want the download error returned", err) + } + + // Neither file may survive: a partial promoted to the final name would be + // indistinguishable from a complete download and would never be repaired. + if _, err := os.Stat(partPath(staging, it)); !os.IsNotExist(err) { + t.Errorf("part file survived a failed download") + } + if _, err := os.Stat(finalPath(staging, it)); !os.IsNotExist(err) { + t.Errorf("a failed download was promoted to the final name") + } +} + +func TestSweepPartialsRemovesOnlyPartFiles(t *testing.T) { + staging := t.TempDir() + keep := filepath.Join(staging, "1234567890_1_done.mp4") + drop := filepath.Join(staging, "1234567890_2_wip.mp4"+partSuffix) + legacy := filepath.Join(staging, "1234567890_3_old.mp4.tmp") + + for _, p := range []string{keep, drop, legacy} { + if err := os.WriteFile(p, []byte("x"), 0o600); err != nil { + t.Fatalf("seed %q: %v", p, err) + } + } + + removed, err := SweepPartials(staging) + if err != nil { + t.Fatalf("SweepPartials: %v", err) + } + if removed != 1 { + t.Errorf("removed = %d, want 1", removed) + } + if _, err := os.Stat(keep); err != nil { + t.Errorf("a completed file was swept: %v", err) + } + // tdl's own suffix is left alone: a shared staging directory during the + // cutover may hold files a legacy run is still writing. + if _, err := os.Stat(legacy); err != nil { + t.Errorf("a legacy tdl .tmp file was swept: %v", err) + } +} + +func TestSweepPartialsOnMissingDirectory(t *testing.T) { + removed, err := SweepPartials(filepath.Join(t.TempDir(), "absent")) + if err != nil { + t.Errorf("SweepPartials on a missing directory = %v, want nil", err) + } + if removed != 0 { + t.Errorf("removed = %d, want 0", removed) + } +} + +// 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) { + staging := t.TempDir() + bad := testItem(t, 7, "../../escape.conf", 10) + + seq := func(yield func(tgsource.Item, error) bool) { yield(bad, nil) } + 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") + } + err := it.Err() + if err == nil { + t.Fatal("Err() = nil after rejecting an unsafe name") + } + if !strings.Contains(err.Error(), "message 7") { + t.Errorf("error should name the message, got: %v", err) + } +} + +func TestElemIterOpensPartFilesAndPropagatesWalkErrors(t *testing.T) { + staging := t.TempDir() + good := testItem(t, 1, "a.mp4", 10) + + t.Run("opens a part file", func(t *testing.T) { + seq := func(yield func(tgsource.Item, error) bool) { yield(good, nil) } + it := newElemIter(seq, staging, true) + defer func() { _ = it.Close() }() + + if !it.Next(t.Context()) { + t.Fatalf("Next() = false, Err() = %v", it.Err()) + } + e := it.Value() + if !e.AsTakeout() { + t.Error("AsTakeout() = false, want the configured value") + } + if e.File().Size() != 10 || e.File().DC() != 2 { + t.Errorf("File() = size %d dc %d, want 10 and 2", e.File().Size(), e.File().DC()) + } + if _, err := os.Stat(partPath(staging, good)); err != nil { + t.Errorf("part file was not created: %v", err) + } + }) + + t.Run("propagates a walk error", func(t *testing.T) { + want := errors.New("history walk failed") + seq := func(yield func(tgsource.Item, error) bool) { yield(tgsource.Item{}, want) } + it := newElemIter(seq, staging, false) + defer func() { _ = it.Close() }() + + if it.Next(t.Context()) { + t.Fatal("Next() = true after a walk error") + } + if !errors.Is(it.Err(), want) { + t.Errorf("Err() = %v, want %v", it.Err(), want) + } + }) +} + +// Byte accounting has to treat ProgressState as a running total, not a delta, +// or the aggregate drifts upward on every callback. +func TestProgressAccountsBytesAsRunningTotals(t *testing.T) { + staging := t.TempDir() + it := testItem(t, 1, "a.mp4", 100) + e := openElem(t, staging, it) + + p := newProgress(func(*elem, error) error { return nil }, nil) + p.OnAdd(e) + p.OnDownload(e, progressState(40)) + p.OnDownload(e, progressState(100)) + p.OnDone(e, nil) + + outcomes, stats := p.results() + if stats.BytesDone != 100 { + t.Errorf("BytesDone = %d, want 100 (states are totals, not deltas)", stats.BytesDone) + } + if stats.BytesTotal != 100 { + t.Errorf("BytesTotal = %d, want 100", stats.BytesTotal) + } + if stats.Done != 1 || stats.Failed != 0 { + t.Errorf("Done/Failed = %d/%d, want 1/0", stats.Done, stats.Failed) + } + if len(outcomes) != 1 || outcomes[0].Err != nil { + t.Errorf("outcomes = %+v, want one success", outcomes) + } +} + +func TestProgressRecordsFailuresWithoutAborting(t *testing.T) { + staging := t.TempDir() + ok := testItem(t, 1, "a.mp4", 10) + bad := testItem(t, 2, "b.mp4", 10) + + p := newProgress(func(*elem, error) error { return nil }, nil) + p.OnDone(openElem(t, staging, ok), nil) + p.OnDone(openElem(t, staging, bad), errors.New("flood wait")) + + outcomes, stats := p.results() + if stats.Done != 1 || stats.Failed != 1 { + t.Errorf("Done/Failed = %d/%d, want 1/1", stats.Done, stats.Failed) + } + if len(outcomes) != 2 { + t.Fatalf("outcomes = %d, want 2 — a failure must be recorded, not dropped", len(outcomes)) + } +} + +func progressState(done int64) downloader.ProgressState { + return downloader.ProgressState{Downloaded: done, Total: 100} +} + +// core's Download logs a failed transfer and returns nil, and OnDone is deferred +// on that named return — so a truncated file reaches finish claiming success. +// The size check is the only thing standing between that and an archived +// fragment, so it is tested directly. +func TestFinishRejectsShortDownloadDespiteNilError(t *testing.T) { + staging := t.TempDir() + it := testItem(t, 3, "video.mp4", 1000) + + e := openElem(t, staging, it) + if _, err := e.file.WriteAt(make([]byte, 400), 0); err != nil { + t.Fatalf("write: %v", err) + } + + // nil, exactly as the downloader reports a failed transfer. + err := finish(staging, e, nil) + if err == nil { + t.Fatal("finish accepted a 400-byte file for a 1000-byte item") + } + if !strings.Contains(err.Error(), "short download") { + t.Errorf("error should name the short download, got: %v", err) + } + if _, err := os.Stat(finalPath(staging, it)); !os.IsNotExist(err) { + t.Error("a truncated file was promoted to its final name") + } + if _, err := os.Stat(partPath(staging, it)); !os.IsNotExist(err) { + t.Error("the truncated part file was left behind") + } +} + +// The size check must not reject a genuinely complete file. +func TestFinishAcceptsExactSize(t *testing.T) { + staging := t.TempDir() + it := testItem(t, 4, "exact.mp4", 2048) + + e := openElem(t, staging, it) + if _, err := e.file.WriteAt(make([]byte, 2048), 0); err != nil { + t.Fatalf("write: %v", err) + } + if err := finish(staging, e, nil); err != nil { + t.Fatalf("finish rejected an exact-size file: %v", err) + } + info, err := os.Stat(finalPath(staging, it)) + if err != nil { + t.Fatalf("final file missing: %v", err) + } + if info.Size() != 2048 { + t.Errorf("final size = %d, want 2048", info.Size()) + } +} diff --git a/internal/pipeline/elem.go b/internal/pipeline/elem.go index 30e39fa..d3c3464 100644 --- a/internal/pipeline/elem.go +++ b/internal/pipeline/elem.go @@ -64,6 +64,12 @@ type elemIter struct { staging string takeout bool + // acquire reserves staging space for the next item. Blocking here is what + // makes backpressure work: core's Download calls Next from its dispatch + // loop, so a blocked Next stops new downloads starting without stopping the + // uploads that free the space. + acquire func(context.Context, int64) error + current *elem err error @@ -104,6 +110,13 @@ func (i *elemIter) Next(ctx context.Context) bool { return false } + if i.acquire != nil { + if err := i.acquire(ctx, item.Size()); err != nil { + i.err = err + return false + } + } + f, err := os.OpenFile(partPath(i.staging, item), os.O_CREATE|os.O_RDWR, 0o600) if err != nil { i.err = fmt.Errorf("open destination for message %d: %w", item.MessageID, err) diff --git a/internal/pipeline/pipeline.go b/internal/pipeline/pipeline.go new file mode 100644 index 0000000..2bcd1b4 --- /dev/null +++ b/internal/pipeline/pipeline.go @@ -0,0 +1,193 @@ +package pipeline + +import ( + "context" + "errors" + "fmt" + "iter" + "sync" + + "github.com/rclone/rclone/fs" + "golang.org/x/sync/semaphore" + + "github.com/iyear/tdl/core/dcpool" + + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +// Options configures a full download-and-upload run. +type Options struct { + Pool dcpool.Pool + Dst fs.Fs + Staging string + + Threads int // connections per file + Limit int // files downloading at once + Uploads int // files uploading at once + Budget int64 // bytes allowed in staging at once; 0 means unbounded + + Confirm bool // re-state each uploaded object to prove its size + Takeout bool + + // MaxFailures trips the run after this many consecutive upload failures. + // Zero uses the shell pipeline's default of 5. + MaxFailures int + + Report func(Stats) +} + +// Result is what a run achieved. +type Result struct { + Stats Stats + Outcomes []Outcome +} + +// Failed lists the items that did not make it to the remote. +func (r Result) Failed() []Outcome { + var out []Outcome + for _, o := range r.Outcomes { + if o.Err != nil { + out = append(out, o) + } + } + return out +} + +// 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 +// it inferred completion from a filename suffix plus a file's age, polled +// `du -sk` every ten seconds, and enforced its disk cap by sending SIGSTOP and +// SIGCONT to the tdl process. None of that exists here. Completion is a function +// returning. The cap is a semaphore: a download acquires its own size before +// starting and releases it only once the upload has confirmed, so when the +// remote is slow the acquire blocks and downloads pause on their own. +// +// Blocking in the iterator is safe by construction — core's Download calls +// Iter.Next from its dispatch loop while workers run in an errgroup, so a +// blocked Next stalls new work without stopping the uploads that free the budget +// (downloader.go:36-63). +func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (Result, error) { + if o.Uploads <= 0 { + o.Uploads = 1 + } + if o.MaxFailures <= 0 { + o.MaxFailures = 5 + } + + local, err := fs.NewFs(ctx, o.Staging) + if err != nil { + return Result{}, fmt.Errorf("open staging directory as a filesystem: %w", err) + } + up := &uploader{local: local, dst: o.Dst, confirm: o.Confirm} + + budget := newBudget(o.Budget) + uploads := make(chan tgsource.Item, o.Uploads) + + // Upload workers own the release side of the budget, so every path out of + // one — success, failure, cancellation — must release, or the run deadlocks + // with downloads waiting on space that is never freed. + var ( + wg sync.WaitGroup + mu sync.Mutex + uploadErrs []error + streak int + tripped bool + ) + upCtx, tripRun := context.WithCancel(ctx) + defer tripRun() + + for range o.Uploads { + wg.Add(1) + go func() { + defer wg.Done() + for it := range uploads { + err := up.upload(upCtx, it) + budget.release(it.Size()) + + mu.Lock() + if err != nil { + uploadErrs = append(uploadErrs, err) + streak++ + if streak >= o.MaxFailures && !tripped { + // A remote that fails this many times running is not + // going to recover on its own, and continuing just fills + // staging until the disk does. + tripped = true + tripRun() + } + } else { + streak = 0 + } + mu.Unlock() + } + }() + } + + dlOutcomes, stats, dlErr := Download(ctx, seq, DownloadOptions{ + Pool: o.Pool, + Staging: o.Staging, + Threads: o.Threads, + Limit: o.Limit, + Takeout: o.Takeout, + Report: o.Report, + acquire: budget.acquire, + onReady: func(it tgsource.Item) { uploads <- it }, + onFailed: func(it tgsource.Item) { + // Nothing was staged, so the reservation has to come back here + // instead of from an upload that will never happen. + budget.release(it.Size()) + }, + }) + + close(uploads) + wg.Wait() + + mu.Lock() + errs := append([]error(nil), uploadErrs...) + trip := tripped + mu.Unlock() + + res := Result{Stats: stats, Outcomes: dlOutcomes} + switch { + case dlErr != nil: + return res, dlErr + case trip: + return res, fmt.Errorf("stopping after %d consecutive upload failures: %w", + o.MaxFailures, errors.Join(errs...)) + case len(errs) > 0: + return res, errors.Join(errs...) + } + return res, nil +} + +// budget bounds how many bytes of downloaded-but-not-yet-uploaded data sit on +// local disk. A zero limit means no bound. +type budget struct{ sem *semaphore.Weighted } + +func newBudget(limit int64) *budget { + if limit <= 0 { + return &budget{} + } + return &budget{sem: semaphore.NewWeighted(limit)} +} + +func (b *budget) acquire(ctx context.Context, n int64) error { + if b.sem == nil { + return nil + } + // An item larger than the whole budget could never be admitted and would + // block forever, so it is refused with an error that says what to change. + // Callers validate up front too; this is the guard for an item whose size + // was not known then. + if err := b.sem.Acquire(ctx, n); err != nil { + return fmt.Errorf("waiting for %d bytes of staging space: %w", n, err) + } + return nil +} + +func (b *budget) release(n int64) { + if b.sem != nil { + b.sem.Release(n) + } +} diff --git a/internal/pipeline/pipeline_test.go b/internal/pipeline/pipeline_test.go new file mode 100644 index 0000000..e36271e --- /dev/null +++ b/internal/pipeline/pipeline_test.go @@ -0,0 +1,147 @@ +package pipeline + +import ( + "context" + "errors" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/iyear/tdl/core/tmedia" + + "github.com/tiennm99dev/telegram-exporter/internal/naming" + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +// The budget is what replaces run.sh's du-polling and SIGSTOP/SIGCONT cap +// draining, so the property it has to hold is simple and worth pinning: the sum +// of outstanding reservations never exceeds the limit. +func TestBudgetBoundsOutstandingBytes(t *testing.T) { + const limit = 1000 + b := newBudget(limit) + ctx := t.Context() + + var ( + mu sync.Mutex + held int64 + peak int64 + wg sync.WaitGroup + acquires atomic.Int64 + ) + + for range 20 { + wg.Add(1) + go func() { + defer wg.Done() + const size = 300 + if err := b.acquire(ctx, size); err != nil { + t.Errorf("acquire: %v", err) + return + } + acquires.Add(1) + + mu.Lock() + held += size + if held > peak { + peak = held + } + mu.Unlock() + + time.Sleep(time.Millisecond) + + mu.Lock() + held -= size + mu.Unlock() + b.release(size) + }() + } + wg.Wait() + + if acquires.Load() != 20 { + t.Errorf("acquired %d times, want 20 — every item must eventually get through", acquires.Load()) + } + if peak > limit { + t.Errorf("peak outstanding = %d bytes, over the %d limit", peak, limit) + } +} + +// A zero limit means the operator asked for no cap; acquiring must not block or +// account, or an unbounded run would stall. +func TestBudgetUnboundedWhenLimitIsZero(t *testing.T) { + b := newBudget(0) + for range 5 { + if err := b.acquire(t.Context(), 1<<40); err != nil { + t.Fatalf("acquire on an unbounded budget: %v", err) + } + } + b.release(1 << 40) // must not panic +} + +// Cancelling must unblock a waiter rather than leaving the run wedged. +func TestBudgetAcquireHonoursCancellation(t *testing.T) { + b := newBudget(100) + if err := b.acquire(t.Context(), 100); err != nil { + t.Fatalf("first acquire: %v", err) + } + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { done <- b.acquire(ctx, 100) }() + + // The second acquire cannot succeed while the first is outstanding. + select { + case err := <-done: + t.Fatalf("acquire succeeded with no space free: %v", err) + case <-time.After(50 * time.Millisecond): + } + + cancel() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) { + t.Errorf("acquire error = %v, want context.Canceled", err) + } + case <-time.After(2 * time.Second): + t.Fatal("cancelling did not unblock the waiter") + } +} + +// An item bigger than the whole budget can never be admitted. It must fail +// rather than hang, because a hang here looks exactly like a slow remote. +func TestBudgetRefusesAnItemLargerThanTheLimit(t *testing.T) { + b := newBudget(100) + ctx, cancel := context.WithTimeout(t.Context(), 500*time.Millisecond) + defer cancel() + + err := b.acquire(ctx, 5000) + if err == nil { + t.Fatal("acquire of an oversized item succeeded") + } + if !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("acquire error = %v, want the deadline to surface", err) + } +} + +func TestResultFailedListsOnlyErrors(t *testing.T) { + r := Result{Outcomes: []Outcome{ + {Item: mustItem(1), Err: nil}, + {Item: mustItem(2), Err: errors.New("flood wait")}, + {Item: mustItem(3), Err: nil}, + {Item: mustItem(4), Err: errors.New("short download")}, + }} + failed := r.Failed() + if len(failed) != 2 { + t.Fatalf("Failed() = %d entries, want 2", len(failed)) + } + for _, f := range failed { + if f.Err == nil { + t.Errorf("Failed() returned a successful outcome: %+v", f) + } + } +} + +func mustItem(id int) tgsource.Item { + m := &tmedia.Media{Name: "f.mp4", Size: 10} + return tgsource.Item{DialogID: 1, MessageID: id, Name: naming.For(1, id, m), Media: m} +} diff --git a/internal/pipeline/progress.go b/internal/pipeline/progress.go new file mode 100644 index 0000000..8c7c75c --- /dev/null +++ b/internal/pipeline/progress.go @@ -0,0 +1,105 @@ +package pipeline + +import ( + "sync" + + "github.com/iyear/tdl/core/downloader" + + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +// Outcome is what happened to one item. +type Outcome struct { + Item tgsource.Item + Err error // nil when the file downloaded and was renamed into place +} + +// Stats is a snapshot of a run's progress. +type Stats struct { + Started int + Done int + Failed int + BytesDone int64 + BytesTotal int64 +} + +// progress collects per-item results and feeds an optional live reporter. +// +// The downloader calls these from its worker goroutines, so everything here is +// mutex-guarded. OnDone fires once per item whether it succeeded or not, which +// is what makes it the right place to finish the file — the downloader itself +// never closes or renames what To() handed it. +type progress struct { + mu sync.Mutex + stats Stats + outcomes []Outcome + inFlight map[int]int64 // message id -> bytes written so far + + finish func(*elem, error) error + report func(Stats) +} + +func newProgress(finish func(*elem, error) error, report func(Stats)) *progress { + return &progress{ + inFlight: make(map[int]int64), + finish: finish, + report: report, + } +} + +func (p *progress) OnAdd(e downloader.Elem) { + el := e.(*elem) + p.mu.Lock() + p.stats.Started++ + p.stats.BytesTotal += el.item.Size() + stats := p.stats + p.mu.Unlock() + p.emit(stats) +} + +func (p *progress) OnDownload(e downloader.Elem, state downloader.ProgressState) { + el := e.(*elem) + p.mu.Lock() + // State carries the running total for this item, not a delta, so the + // aggregate is adjusted by the difference since the last callback. + prev := p.inFlight[el.item.MessageID] + p.inFlight[el.item.MessageID] = state.Downloaded + p.stats.BytesDone += state.Downloaded - prev + stats := p.stats + p.mu.Unlock() + p.emit(stats) +} + +func (p *progress) OnDone(e downloader.Elem, err error) { + el := e.(*elem) + + // Closing and renaming happens here because this is the only callback that + // runs exactly once per item and knows whether it succeeded. + if ferr := p.finish(el, err); ferr != nil && err == nil { + err = ferr + } + + p.mu.Lock() + delete(p.inFlight, el.item.MessageID) + if err != nil { + p.stats.Failed++ + } else { + p.stats.Done++ + } + p.outcomes = append(p.outcomes, Outcome{Item: el.item, Err: err}) + stats := p.stats + p.mu.Unlock() + p.emit(stats) +} + +func (p *progress) emit(s Stats) { + if p.report != nil { + p.report(s) + } +} + +func (p *progress) results() ([]Outcome, Stats) { + p.mu.Lock() + defer p.mu.Unlock() + return p.outcomes, p.stats +} diff --git a/internal/pipeline/upload.go b/internal/pipeline/upload.go new file mode 100644 index 0000000..89ed3eb --- /dev/null +++ b/internal/pipeline/upload.go @@ -0,0 +1,47 @@ +package pipeline + +import ( + "context" + "fmt" + + "github.com/rclone/rclone/fs" + "github.com/rclone/rclone/fs/operations" + + "github.com/tiennm99dev/telegram-exporter/internal/tgsource" +) + +// uploader moves finished files from staging to the destination remote. +type uploader struct { + local fs.Fs // the staging directory as an rclone filesystem + dst fs.Fs + confirm bool +} + +// upload moves one finished file to the remote and, unless disabled, proves it +// arrived at the expected size. +// +// 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 — which is what lets the +// byte budget be released. +// +// Confirmation closes a gap the shell pipeline left open: there, a truncated +// upload was only noticed by a later verify pass, after the local copy was +// already gone. Re-stating the object costs one round trip per file and turns a +// silent corruption into a retry. +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) + } + if !u.confirm { + return nil + } + + obj, err := u.dst.NewObject(ctx, it.Name) + if err != nil { + return fmt.Errorf("confirm %q: %w", it.Name, err) + } + if got := obj.Size(); got != it.Size() { + return fmt.Errorf("confirm %q: remote has %d bytes, expected %d", it.Name, got, it.Size()) + } + return nil +} diff --git a/internal/remote/fs.go b/internal/remote/fs.go index 6e71596..1ab57b0 100644 --- a/internal/remote/fs.go +++ b/internal/remote/fs.go @@ -98,6 +98,20 @@ func Resolve(ctx context.Context, remote string) (fs.Fs, error) { return f, nil } +// EnsureDir creates the destination if it is not there yet. +// +// A destination that does not exist yet is the normal case for a first run, and +// the shell pipeline created it up front for the same reason — the call doubles +// as the reachability and credentials check, since a remote that refuses a +// mkdir will refuse the uploads too. Doing it before the chat is read means a +// bad destination fails in seconds rather than after a full history walk. +func EnsureDir(ctx context.Context, f fs.Fs) error { + if err := f.Mkdir(ctx, ""); err != nil { + return fmt.Errorf("cannot create %q — check credentials and connectivity: %w", f.String(), err) + } + return nil +} + // FreeBytes reports free space on the remote. // // Backends without quota reporting return ok=false rather than an error: the diff --git a/internal/remote/index.go b/internal/remote/index.go index 73bd950..40acc9c 100644 --- a/internal/remote/index.go +++ b/internal/remote/index.go @@ -2,6 +2,7 @@ package remote import ( "context" + "errors" "fmt" "path" "slices" @@ -88,6 +89,12 @@ func BuildIndex(ctx context.Context, f fs.Fs, dialogID int64) (*Index, error) { idx.byID[id] = append(idx.byID[id], name) } }); err != nil { + // A destination that does not exist yet holds nothing. That is an empty + // index, not a failure — it is what a first run against a new path looks + // like, and treating it as an error would make verify unusable there. + if errors.Is(err, fs.ErrorDirNotFound) { + return idx, nil + } return nil, fmt.Errorf("list %s: %w", f.String(), err) } return idx, nil diff --git a/internal/report/progress.go b/internal/report/progress.go new file mode 100644 index 0000000..e458fa8 --- /dev/null +++ b/internal/report/progress.go @@ -0,0 +1,108 @@ +// Package report renders a run's progress for whoever is watching. +package report + +import ( + "fmt" + "io" + "os" + "sync" + "time" + + "github.com/tiennm99dev/telegram-exporter/internal/pipeline" +) + +// statsInterval is how often a redirected run prints a line. +// +// The shell pipeline learned this the hard way: tdl's progress bar is ANSI +// redraws, which are right on a terminal and turn a captured log into +// megabytes of control characters. So a TTY gets a redrawn line and everything +// else gets a periodic summary. +const statsInterval = 30 * time.Second + +// Reporter renders progress, adapting to whether it is writing to a terminal. +type Reporter struct { + w io.Writer + tty bool + total int + totalBytes int64 + + mu sync.Mutex + started time.Time + lastLine time.Time +} + +// New builds a reporter for w, which is treated as a terminal when it is one. +// +// The totals are passed in rather than taken from Stats because Stats.BytesTotal +// only counts items the downloader has started, so a progress line built from it +// shows a denominator that grows as the run proceeds — "0 B of 52 KiB" on a run +// that will move gigabytes. The caller knows the real figures before starting. +func New(w io.Writer, total int, totalBytes int64) *Reporter { + return &Reporter{w: w, tty: isTerminal(w), total: total, totalBytes: totalBytes, started: time.Now()} +} + +// Update renders a snapshot. Safe to call from several goroutines, and cheap +// enough to call on every progress callback. +func (r *Reporter) Update(s pipeline.Stats) { + r.mu.Lock() + defer r.mu.Unlock() + + now := time.Now() + if r.tty { + // \r rather than \n: one line, redrawn. + fmt.Fprintf(r.w, "\r\033[K%s", r.line(s, now)) + return + } + if now.Sub(r.lastLine) < statsInterval { + return + } + r.lastLine = now + fmt.Fprintf(r.w, "%s\n", r.line(s, now)) +} + +// Finish writes the closing summary, ending the redrawn line if there was one. +func (r *Reporter) Finish(s pipeline.Stats) { + r.mu.Lock() + defer r.mu.Unlock() + if r.tty { + fmt.Fprint(r.w, "\r\033[K") + } + elapsed := time.Since(r.started).Round(time.Second) + fmt.Fprintf(r.w, "%d done, %d failed, %s in %s (%s/s)\n", + s.Done, s.Failed, humanBytes(s.BytesDone), elapsed, + humanBytes(int64(float64(s.BytesDone)/max(elapsed.Seconds(), 1)))) +} + +func (r *Reporter) line(s pipeline.Stats, now time.Time) string { + elapsed := now.Sub(r.started) + rate := float64(s.BytesDone) / max(elapsed.Seconds(), 1) + return fmt.Sprintf("%d/%d done, %d failed, %s of %s, %s/s", + s.Done, r.total, s.Failed, + humanBytes(s.BytesDone), humanBytes(r.totalBytes), humanBytes(int64(rate))) +} + +func humanBytes(n int64) string { + const unit = 1024 + if n < unit { + return fmt.Sprintf("%d B", n) + } + div, exp := int64(unit), 0 + for v := n / unit; v >= unit; v /= unit { + div *= unit + exp++ + } + return fmt.Sprintf("%.1f %ciB", float64(n)/float64(div), "KMGT"[exp]) +} + +// isTerminal reports whether w is a character device. +func isTerminal(w io.Writer) bool { + f, ok := w.(*os.File) + if !ok { + return false + } + info, err := f.Stat() + if err != nil { + return false + } + return info.Mode()&os.ModeCharDevice != 0 +}