From fafeab47500b40ad26f9dc9bd48ea31befb80281 Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Sun, 6 Sep 2026 21:05:59 +0700 Subject: [PATCH] feat: report progress while reading a chat and listing a remote MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both phases ran for minutes printing nothing between their opening line and their result, so a working run looked exactly like a hung one — which is how it was reported. The message counter lives in Walk rather than in the caller's loop because most of a chat is not media: text-only and service messages are filtered out inside the walk, so a caller counting yielded items still sees nothing while crossing a long stretch of conversation. Cadence follows the existing reporter: a terminal redraws one line, a redirected run gets a periodic one, since ANSI redraws are what turned the shell pipeline's captured logs into megabytes of control characters. --- cmd/tgexport/list.go | 4 +- cmd/tgexport/sync.go | 28 ++++++++-- cmd/tgexport/verify.go | 34 ++++++++---- internal/remote/index.go | 9 +++- internal/remote/index_env_test.go | 2 +- internal/remote/index_test.go | 2 +- internal/report/ticker.go | 89 +++++++++++++++++++++++++++++++ internal/report/ticker_test.go | 59 ++++++++++++++++++++ internal/tgsource/iterate.go | 14 ++++- 9 files changed, 222 insertions(+), 19 deletions(-) create mode 100644 internal/report/ticker.go create mode 100644 internal/report/ticker_test.go diff --git a/cmd/tgexport/list.go b/cmd/tgexport/list.go index 81981af..5cdccb0 100644 --- a/cmd/tgexport/list.go +++ b/cmd/tgexport/list.go @@ -11,6 +11,7 @@ import ( "github.com/iyear/tdl/core/dcpool" tdlstorage "github.com/iyear/tdl/core/storage" + "github.com/tiennm99dev/telegram-exporter/internal/report" "github.com/tiennm99dev/telegram-exporter/internal/tdlkv" "github.com/tiennm99dev/telegram-exporter/internal/tgsource" ) @@ -66,9 +67,10 @@ func listCmd(ctx context.Context, args []string) error { // swallowed by a deferred call nobody checks. out := bufio.NewWriter(os.Stdout) + scan := report.NewTicker(os.Stderr, "messages read") count := 0 var total int64 - for it, err := range tgsource.Walk(ctx, api, peer) { + for it, err := range tgsource.Walk(ctx, api, peer, scan.Update) { if err != nil { return err } diff --git a/cmd/tgexport/sync.go b/cmd/tgexport/sync.go index b87f6c1..6b3a6a1 100644 --- a/cmd/tgexport/sync.go +++ b/cmd/tgexport/sync.go @@ -122,15 +122,21 @@ func syncCmd(ctx context.Context, args []string) error { } fmt.Fprintf(os.Stderr, "reading %s\n", *chat) + scan := report.NewTicker(os.Stderr, "messages read") var items []tgsource.Item - for it, err := range tgsource.Walk(ctx, api, peer) { + var scanned int + for it, err := range tgsource.Walk(ctx, api, peer, func(n int) { + scanned = n + scan.Update(n) + }) { if err != nil { return err } items = append(items, it) } + scan.Done(scanned) - idx, err := remote.BuildIndex(ctx, dst, peer.ID()) + idx, err := indexRemote(ctx, dst, peer.ID()) if err != nil { return err } @@ -191,7 +197,7 @@ func syncCmd(ctx context.Context, args []string) error { // 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()) + idx, err = indexRemote(ctx, dst, peer.ID()) if err != nil { return err } @@ -280,6 +286,22 @@ func validateBudget(budget int64, todo []tgsource.Item) error { return nil } +// indexRemote lists the destination, reporting progress as it goes. +func indexRemote(ctx context.Context, dst fs.Fs, dialogID int64) (*remote.Index, error) { + fmt.Fprintf(os.Stderr, "indexing %s\n", dst.String()) + tick := report.NewTicker(os.Stderr, "objects listed") + var seen int + idx, err := remote.BuildIndex(ctx, dst, dialogID, func(n int) { + seen = n + tick.Update(n) + }) + if err != nil { + return nil, err + } + tick.Done(seen) + return idx, nil +} + // warnCollisions reports basenames the index found at more than one path. // // verify printed this and sync did not, which was backwards: an ambiguous diff --git a/cmd/tgexport/verify.go b/cmd/tgexport/verify.go index d79d1e7..8239e47 100644 --- a/cmd/tgexport/verify.go +++ b/cmd/tgexport/verify.go @@ -15,6 +15,7 @@ import ( "github.com/rclone/rclone/fs/operations" "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" @@ -66,7 +67,7 @@ func verifyCmd(ctx context.Context, args []string) error { } var ( - report verify.Report + result verify.Report idx *remote.Index ) if err := sess.Run(ctx, func(ctx context.Context, pool dcpool.Pool) error { @@ -80,19 +81,30 @@ func verifyCmd(ctx context.Context, args []string) error { // The chat is walked first and the remote listed second, so the snapshot // is never older than the wanted set. The reverse order could report a // file absent that was uploaded while the walk was still running. + fmt.Fprintf(os.Stderr, "reading %s\n", *chat) + scan := report.NewTicker(os.Stderr, "messages read") + var scanned int var items []tgsource.Item - for it, err := range tgsource.Walk(ctx, api, peer) { + for it, err := range tgsource.Walk(ctx, api, peer, func(n int) { scanned = n; scan.Update(n) }) { if err != nil { return err } items = append(items, it) } - idx, err = remote.BuildIndex(ctx, dst, peer.ID()) + scan.Done(scanned) + + fmt.Fprintf(os.Stderr, "indexing %s\n", dst.String()) + idxTick := report.NewTicker(os.Stderr, "objects listed") + var listed int + idx, err = remote.BuildIndex(ctx, dst, peer.ID(), func(n int) { + listed = n + idxTick.Update(n) + }) if err != nil { return err } - fmt.Fprintf(os.Stderr, "indexed %d objects on %s\n", idx.Len(), dst.String()) + idxTick.Done(listed) if dup := idx.Collisions(); len(dup) > 0 { // An ambiguous snapshot makes every verdict about these names // unreliable, so it is reported rather than silently resolved. @@ -104,26 +116,26 @@ func verifyCmd(ctx context.Context, args []string) error { } fmt.Fprintln(os.Stderr) - report = verify.Check(items, idx) + result = verify.Check(items, idx) return nil }); err != nil { return err } out := bufio.NewWriter(os.Stdout) - report.Write(out) + result.Write(out) if err := out.Flush(); err != nil { return err } if *delStale { - if err := deleteMisnamed(ctx, dst, idx, report, *assumeYes); err != nil { + if err := deleteMisnamed(ctx, dst, idx, result, *assumeYes); err != nil { return err } } - if !report.Complete() { - return fmt.Errorf("%w: %d file(s) still to fetch", errIncomplete, len(report.Todo())) + if !result.Complete() { + return fmt.Errorf("%w: %d file(s) still to fetch", errIncomplete, len(result.Todo())) } return nil } @@ -138,9 +150,9 @@ func verifyCmd(ctx context.Context, args []string) error { // used elsewhere. Those differ the moment a remote has directory structure, and // deleting by basename would either miss the object or — worse, if the root // happens to hold a same-named file — delete the wrong one. -func deleteMisnamed(ctx context.Context, dst rclonefs.Fs, idx *remote.Index, report verify.Report, assumeYes bool) error { +func deleteMisnamed(ctx context.Context, dst rclonefs.Fs, idx *remote.Index, result verify.Report, assumeYes bool) error { var targets []string - for _, m := range report.Misnamed { + for _, m := range result.Misnamed { targets = append(targets, m.Found...) } if len(targets) == 0 { diff --git a/internal/remote/index.go b/internal/remote/index.go index 4e97e5c..c95af98 100644 --- a/internal/remote/index.go +++ b/internal/remote/index.go @@ -62,7 +62,9 @@ type Index struct { // why every field that can narrow a listing is set explicitly. A zero-value // Options is not a substitute either — it fails validation, because MinAge and // MaxAge both being 0 reads as "min > max". -func BuildIndex(ctx context.Context, f fs.Fs, dialogID int64) (*Index, error) { +// onCount, when non-nil, is called with the number of objects seen so far. +// Listing a remote of any size takes minutes and says nothing while it runs. +func BuildIndex(ctx context.Context, f fs.Fs, dialogID int64, onCount func(n int)) (*Index, error) { ctx, ci := fs.AddConfig(ctx) ci.MaxDepth = -1 @@ -88,7 +90,12 @@ func BuildIndex(ctx context.Context, f fs.Fs, dialogID int64) (*Index, error) { } // ListFn is documented not to call fn concurrently, so the maps need no lock. + seen := 0 if err := operations.ListFn(ctx, f, func(o fs.Object) { + seen++ + if onCount != nil { + onCount(seen) + } name := path.Base(o.Remote()) if _, seen := idx.byName[name]; seen { diff --git a/internal/remote/index_env_test.go b/internal/remote/index_env_test.go index 1fb4124..d4a05a3 100644 --- a/internal/remote/index_env_test.go +++ b/internal/remote/index_env_test.go @@ -52,7 +52,7 @@ func indexChild(t *testing.T) { if err != nil { t.Fatal(err) } - idx, err := BuildIndex(t.Context(), f, 1234567890) + idx, err := BuildIndex(t.Context(), f, 1234567890, nil) if err != nil { t.Fatalf("BuildIndex: %v", err) } diff --git a/internal/remote/index_test.go b/internal/remote/index_test.go index 5abd8c4..eeb1b15 100644 --- a/internal/remote/index_test.go +++ b/internal/remote/index_test.go @@ -33,7 +33,7 @@ func localIndex(t *testing.T, files map[string]int) *Index { if err != nil { t.Fatalf("open local fs: %v", err) } - idx, err := BuildIndex(ctx, f, testDialog) + idx, err := BuildIndex(ctx, f, testDialog, nil) if err != nil { t.Fatalf("BuildIndex: %v", err) } diff --git a/internal/report/ticker.go b/internal/report/ticker.go new file mode 100644 index 0000000..c674521 --- /dev/null +++ b/internal/report/ticker.go @@ -0,0 +1,89 @@ +package report + +import ( + "fmt" + "io" + "sync" + "time" +) + +// tickerInterval is how often a terminal redraws a phase counter. Fast enough +// to look alive, slow enough not to matter. +const tickerInterval = 250 * time.Millisecond + +// Ticker reports progress through a long phase that would otherwise be silent. +// +// Reading a chat's history and listing a remote each take minutes on an archive +// of any size, and both used to print nothing between their opening line and +// their result. A run that is working looked identical to one that had hung, so +// the only way to tell was to wait it out. +// +// Cadence follows the same rule as Reporter: a terminal gets a redrawn line, a +// redirected run gets a periodic one, because ANSI redraws turn a captured log +// into megabytes of control characters. +type Ticker struct { + w io.Writer + tty bool + noun string + start time.Time + + mu sync.Mutex + lastLine time.Time +} + +// NewTicker builds a ticker that counts noun, e.g. "messages" or "objects". +func NewTicker(w io.Writer, noun string) *Ticker { + return &Ticker{w: w, tty: isTerminal(w), noun: noun, start: time.Now()} +} + +// Update reports a running count. Safe to call from several goroutines, and +// cheap enough to call per item. +func (t *Ticker) Update(n int) { + if !t.mu.TryLock() { + return + } + defer t.mu.Unlock() + + now := time.Now() + interval := statsInterval + if t.tty { + interval = tickerInterval + } + if now.Sub(t.lastLine) < interval { + return + } + t.lastLine = now + + if t.tty { + fmt.Fprintf(t.w, "\r\033[K %s %s...", humanCount(n), t.noun) + return + } + fmt.Fprintf(t.w, " %s %s...\n", humanCount(n), t.noun) +} + +// Done clears the redrawn line and states the final count. +func (t *Ticker) Done(n int) { + t.mu.Lock() + defer t.mu.Unlock() + if t.tty { + fmt.Fprint(t.w, "\r\033[K") + } + fmt.Fprintf(t.w, " %s %s in %s\n", humanCount(n), t.noun, + time.Since(t.start).Round(time.Second)) +} + +// humanCount groups thousands, so 12000 reads as 12,000. +func humanCount(n int) string { + s := fmt.Sprintf("%d", n) + if len(s) <= 3 { + return s + } + out := make([]byte, 0, len(s)+len(s)/3) + for i, c := range []byte(s) { + if i > 0 && (len(s)-i)%3 == 0 { + out = append(out, ',') + } + out = append(out, c) + } + return string(out) +} diff --git a/internal/report/ticker_test.go b/internal/report/ticker_test.go new file mode 100644 index 0000000..66b3ffa --- /dev/null +++ b/internal/report/ticker_test.go @@ -0,0 +1,59 @@ +package report + +import ( + "strings" + "testing" + "time" +) + +func TestHumanCountGroupsThousands(t *testing.T) { + cases := map[int]string{ + 0: "0", 7: "7", 999: "999", 1000: "1,000", + 12000: "12,000", 11406: "11,406", 1234567: "1,234,567", + } + for in, want := range cases { + if got := humanCount(in); got != want { + t.Errorf("humanCount(%d) = %q, want %q", in, got, want) + } + } +} + +// A redirected run must not accumulate ANSI redraws: that is what turned the +// shell pipeline's captured logs into megabytes of control characters. +func TestTickerWritesNoAnsiWhenRedirected(t *testing.T) { + var sb strings.Builder + tick := NewTicker(&sb, "messages read") + + for i := 1; i <= 5000; i++ { + tick.Update(i) + } + tick.Done(5000) + + out := sb.String() + if strings.ContainsAny(out, "\r\033") { + t.Errorf("redirected output contains control characters: %q", out) + } + if !strings.Contains(out, "5,000 messages read") { + t.Errorf("final count missing from %q", out) + } +} + +// Update is called once per message on an 18k-message walk, so it has to be +// cheap: at most one line per interval, however often it is called. +func TestTickerThrottlesUpdates(t *testing.T) { + var sb strings.Builder + tick := NewTicker(&sb, "objects listed") + tick.lastLine = time.Now() // inside the interval from the start + + for i := 1; i <= 10000; i++ { + tick.Update(i) + } + if n := strings.Count(sb.String(), "\n"); n != 0 { + t.Errorf("wrote %d lines inside one interval, want 0", n) + } + + tick.Done(10000) + if !strings.Contains(sb.String(), "10,000 objects listed in") { + t.Errorf("Done did not state the total: %q", sb.String()) + } +} diff --git a/internal/tgsource/iterate.go b/internal/tgsource/iterate.go index 87c90db..ae44dbc 100644 --- a/internal/tgsource/iterate.go +++ b/internal/tgsource/iterate.go @@ -47,12 +47,24 @@ func (i Item) Size() int64 { return i.Media.Size } // the export JSON encoded as an empty "file" field. On error the sequence yields // a zero Item with that error and stops; cancelling ctx stops it too, so an // interrupted run does not keep paging. -func Walk(ctx context.Context, api *tg.Client, peer peers.Peer) iter.Seq2[Item, error] { +// onScan, when non-nil, is called with the number of messages read so far. It +// has to live here rather than in the caller's loop because most of a chat is +// not media: text-only and service messages are filtered out below, so a caller +// counting yielded items sees nothing at all while the walk crosses a long +// stretch of conversation, and a working run is indistinguishable from a hung +// one. +func Walk(ctx context.Context, api *tg.Client, peer peers.Peer, onScan func(scanned int)) iter.Seq2[Item, error] { return func(yield func(Item, error) bool) { dialogID := peer.ID() it := query.NewQuery(api).Messages().GetHistory(peer.InputPeer()).BatchSize(100).Iter() + scanned := 0 for it.Next(ctx) { + scanned++ + if onScan != nil { + onScan(scanned) + } + msg, ok := it.Value().Msg.(*tg.Message) if !ok { continue // service messages have no media