From 846c2bafe972bb92ab907cf0f8cb7c06552a5357 Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Mon, 7 Sep 2026 00:05:39 +0700 Subject: [PATCH] feat: count download and upload separately MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One combined figure could say how much had moved but not which half was moving it. That is the question a slowing run actually raises: is Telegram the bottleneck, or the remote? The two legs now have their own totals, files and bytes each. The gap between them is the useful part. They normally track a file or two apart; a widening gap is the remote falling behind, which is also staging filling up — visible now before the byte cap starts throttling downloads. Uploads advance a whole file at a time because rclone's MoveFile is one blocking call with no byte callbacks, so a partial upload counts for nothing until it lands. That is the honest reading anyway: what the upload total reports is what is actually on the remote. Both renderers read the same counters, so the terminal and a captured log cannot disagree about the numbers. --- README.md | 20 ++++++++--- internal/report/live.go | 68 ++++++++++++++++++++++--------------- internal/report/progress.go | 36 +++++++++++++++----- internal/report/totals.go | 54 +++++++++++++++++++++++++++++ 4 files changed, 137 insertions(+), 41 deletions(-) create mode 100644 internal/report/totals.go diff --git a/README.md b/README.md index f7d95fc..77a4df8 100644 --- a/README.md +++ b/README.md @@ -76,13 +76,23 @@ indexing PikPak root 'mychannel' staging ./staging, capped at 40.0 GiB concurrency 2 download(s) x 4 thread(s), 2 upload(s) - total 26/606 files [=> ] 617.5 MiB / 79.0 GiB 2.5 MiB/s 8h47m - ↓ …3214_4242_1000000000000000001.mp4 [=======> ] 41.2 MiB / 96.0 MiB 1.8 MiB/s - ↑ …3214_4243_1000000000000000002.mp4 ⠹ uploading 1.9 GiB + ↓ total 26/606 files [=> ] 617.5 MiB / 79.0 GiB 2.5 MiB/s 8h47m + ↑ total 24/606 files [=> ] 598.0 MiB / 79.0 GiB 2.4 MiB/s 8h58m + ↓ …3214_4242_1000000000000000001.mp4 [=======> ] 41.2 MiB / 96.0 MiB 1.8 MiB/s + ↑ …3214_4243_1000000000000000002.mp4 ⠹ uploading 1.9 GiB ``` -Redirected output gets the same information as plain periodic lines plus one -line per archived file, with no cursor movement — a captured log stays readable. +The two legs are counted separately because they run at different speeds and +fail for different reasons. They normally track a file or two apart; a widening +gap means the remote is falling behind and staging is filling up. + +Redirected output gets the same two figures as plain periodic lines, plus one +line per archived file, with no cursor movement — a captured log stays readable: + +``` + download 1,200/606 files, 45.0 GiB of 79.0 GiB, 76.8 MiB/s, ETA 7m33s + upload 1,190/606 files, 44.2 GiB of 79.0 GiB, 75.4 MiB/s, ETA 7m53s +``` `CHAT` accepts a numeric id as printed by `tdl chat ls`, a username with or without `@`, or a `t.me`/`tg://` link. A Bot API `-100…` id is converted diff --git a/internal/report/live.go b/internal/report/live.go index 094f30b..ec68f64 100644 --- a/internal/report/live.go +++ b/internal/report/live.go @@ -17,7 +17,7 @@ import ( // to 100+ characters and the bar has to fit beside them. const nameWidth = 34 -// Live renders a run as a set of progress bars: one overall, plus one for each +// Live renders a run as a set of progress bars: one per leg, plus one for each // file currently moving. // // The aggregate line it replaces could say how much was done but never what was @@ -25,20 +25,27 @@ const nameWidth = 34 // a slow upload, how long the rest would take. On a run measured in hours those // are the only questions worth answering. // +// Download and upload are counted separately because they run at different +// speeds and fail for different reasons. A single combined figure hides the one +// thing worth knowing when a run slows down: whether Telegram or the remote is +// the bottleneck. The two normally track each other a file or two apart; a +// widening gap is the remote falling behind, and staging filling up. +// // Bars are for terminals only. A redirected run gets lineEvents instead, because // this writes ANSI cursor movement continuously and a captured log of it is // unreadable. type Live struct { - w io.Writer - p *mpb.Progress - total *mpb.Bar + w io.Writer + p *mpb.Progress + dlTotal *mpb.Bar + upTotal *mpb.Bar mu sync.Mutex down map[int]*mpb.Bar up map[int]*mpb.Bar - // done is read by the overall bar's decorator from mpb's render goroutine, - // so it is guarded by the same lock as the maps. - done int + // The leg counters are read by the total bars' decorators from mpb's render + // goroutine, so they are guarded by the same lock as the maps. + legs legTotals // seq gives each per-file bar a distinct, increasing priority so bars keep // their position between frames. Sharing one priority lets mpb reorder them @@ -71,13 +78,22 @@ func NewLive(w io.Writer, files int, bytes int64) *Live { files: files, start: time.Now(), } - l.total = p.New(bytes, + l.dlTotal = l.leg(bytes, 0, " ↓ total", func() int { return l.legs.downCount() }) + l.upTotal = l.leg(bytes, 1, " ↑ total", func() int { return l.legs.upCount() }) + return l +} + +// leg builds one of the two whole-run bars. +func (l *Live) leg(bytes int64, priority int, label string, count func() int) *mpb.Bar { + return l.p.New(bytes, mpb.BarStyle().Lbound("[").Filler("=").Tip(">").Padding(" ").Rbound("]"), - mpb.BarPriority(0), + mpb.BarPriority(priority), mpb.BarNoPop(), mpb.PrependDecorators( - decor.Name(" total ", decor.WC{W: 9}), - decor.Any(func(decor.Statistics) string { return l.counts() }, decor.WC{W: 16}), + decor.Name(label+" ", decor.WC{W: 11}), + decor.Any(func(decor.Statistics) string { + return fmt.Sprintf("%s/%s files", humanCount(count()), humanCount(l.files)) + }, decor.WC{W: 16}), ), mpb.AppendDecorators( decor.CountersKibiByte("% .1f / % .1f", decor.WC{W: 20}), @@ -87,11 +103,10 @@ func NewLive(w io.Writer, files int, bytes int64) *Live { decor.OnComplete(decor.AverageETA(decor.ET_STYLE_GO, decor.WC{W: 13}), ""), ), ) - return l } -// Bars are grouped by band: the total on top, then downloads, then uploads. -// Within a band they are ordered by when they started. +// Bars are grouped by band: the two totals on top, then per-file downloads, +// then per-file uploads. Within a band they are ordered by when they started. const ( downloadBand = 1 << 20 uploadBand = 1 << 21 @@ -105,20 +120,12 @@ func (l *Live) next(band int) int { return band + l.seq } -func (l *Live) counts() string { - l.mu.Lock() - defer l.mu.Unlock() - return fmt.Sprintf("%s/%s files", humanCount(l.done), humanCount(l.files)) -} - -// Stats advances the overall bar. +// Stats advances the download bar. func (l *Live) Stats(s pipeline.Stats) { - l.mu.Lock() - l.done = s.Done + s.Failed - l.mu.Unlock() + l.legs.setDownload(s.Done+s.Failed, s.BytesDone) // The bar's own counters, speed and ETA all derive from this, so it is what - // makes the totals move rather than sitting at zero. - l.total.SetCurrent(s.BytesDone) + // makes the total move rather than sitting at zero. + l.dlTotal.SetCurrent(s.BytesDone) } func (l *Live) DownloadStart(it tgsource.Item) { @@ -189,7 +196,7 @@ func (l *Live) UploadStart(it tgsource.Item) { l.mu.Unlock() } -func (l *Live) UploadDone(it tgsource.Item, _ error) { +func (l *Live) UploadDone(it tgsource.Item, err error) { l.mu.Lock() bar := l.up[it.MessageID] delete(l.up, it.MessageID) @@ -197,6 +204,10 @@ func (l *Live) UploadDone(it tgsource.Item, _ error) { if bar != nil { bar.Abort(true) } + if err != nil { + return // nothing reached the remote, so the upload total does not move + } + l.upTotal.SetCurrent(l.legs.addUpload(it.Size())) } // Finish drains the bars and prints the closing summary. @@ -212,7 +223,8 @@ func (l *Live) Finish(s pipeline.Stats) { clear(l.up) l.mu.Unlock() - l.total.Abort(true) + l.dlTotal.Abort(true) + l.upTotal.Abort(true) l.p.Wait() writeSummary(l.w, s, time.Since(l.start)) } diff --git a/internal/report/progress.go b/internal/report/progress.go index 4e73b43..ee726ee 100644 --- a/internal/report/progress.go +++ b/internal/report/progress.go @@ -31,6 +31,8 @@ type Reporter struct { total int totalBytes int64 + legs legTotals + mu sync.Mutex started time.Time lastLine time.Time @@ -64,6 +66,8 @@ func newReporter(w io.Writer, total int, totalBytes int64) *Reporter { // make write latency throttle the downloads themselves. A skipped frame costs // nothing; the next callback is milliseconds away and Finish always prints. func (r *Reporter) Stats(s pipeline.Stats) { + r.legs.setDownload(s.Done+s.Failed, s.BytesDone) + if !r.mu.TryLock() { return } @@ -74,7 +78,9 @@ func (r *Reporter) Stats(s pipeline.Stats) { return } r.lastLine = now - fmt.Fprintf(r.w, "%s\n", r.line(s, now)) + for _, line := range r.lines(now) { + fmt.Fprintf(r.w, "%s\n", line) + } } // The per-file events are recorded as one line each rather than a bar. At one @@ -93,6 +99,9 @@ func (r *Reporter) DownloadDone(it tgsource.Item, err error) { } func (r *Reporter) UploadDone(it tgsource.Item, err error) { + if err == nil { + r.legs.addUpload(it.Size()) + } r.mu.Lock() defer r.mu.Unlock() if err != nil { @@ -109,17 +118,28 @@ func (r *Reporter) Finish(s pipeline.Stats) { writeSummary(r.w, s, time.Since(r.started)) } -func (r *Reporter) line(s pipeline.Stats, now time.Time) string { +// lines renders one line per leg. Separately, because the two run at different +// speeds and a single combined figure hides which of them is the bottleneck — +// the question actually being asked when a run slows down. +func (r *Reporter) lines(now time.Time) []string { elapsed := now.Sub(r.started) - rate := float64(s.BytesDone) / max(elapsed.Seconds(), 1) + dlFiles, dlBytes, upFiles, upBytes := r.legs.snapshot() + return []string{ + r.leg("download", dlFiles, dlBytes, elapsed), + r.leg("upload ", upFiles, upBytes, elapsed), + } +} + +func (r *Reporter) leg(label string, files int, bytes int64, elapsed time.Duration) string { + rate := float64(bytes) / max(elapsed.Seconds(), 1) eta := "—" - if rate > 0 && r.totalBytes > s.BytesDone { - eta = time.Duration(float64(r.totalBytes-s.BytesDone) / rate * float64(time.Second)). + if rate > 0 && r.totalBytes > bytes { + eta = time.Duration(float64(r.totalBytes-bytes) / rate * float64(time.Second)). Round(time.Second).String() } - return fmt.Sprintf(" %s/%s files, %d failed, %s of %s, %s/s, ETA %s", - humanCount(s.Done), humanCount(r.total), s.Failed, - humanBytes(s.BytesDone), humanBytes(r.totalBytes), humanBytes(int64(rate)), eta) + return fmt.Sprintf(" %s %s/%s files, %s of %s, %s/s, ETA %s", + label, humanCount(files), humanCount(r.total), + humanBytes(bytes), humanBytes(r.totalBytes), humanBytes(int64(rate)), eta) } func humanBytes(n int64) string { diff --git a/internal/report/totals.go b/internal/report/totals.go new file mode 100644 index 0000000..1db72ca --- /dev/null +++ b/internal/report/totals.go @@ -0,0 +1,54 @@ +package report + +import "sync" + +// legTotals counts what each half of the run has moved. +// +// The two legs are tracked separately because they report differently: the +// downloader emits byte-level progress, while an upload is one blocking rclone +// call that only reports on completion. Keeping the counters here rather than in +// each renderer means the terminal and the log agree on the same numbers. +type legTotals struct { + mu sync.Mutex + dlFiles int + dlBytes int64 + upFiles int + upBytes int64 +} + +// setDownload records the downloader's running totals. +func (t *legTotals) setDownload(files int, bytes int64) { + t.mu.Lock() + defer t.mu.Unlock() + t.dlFiles, t.dlBytes = files, bytes +} + +// addUpload records one file that reached the remote and returns the new total. +// A whole file at a time is all that is knowable: rclone's MoveFile does not +// report progress within a transfer. +func (t *legTotals) addUpload(size int64) int64 { + t.mu.Lock() + defer t.mu.Unlock() + t.upFiles++ + t.upBytes += size + return t.upBytes +} + +func (t *legTotals) downCount() int { + t.mu.Lock() + defer t.mu.Unlock() + return t.dlFiles +} + +func (t *legTotals) upCount() int { + t.mu.Lock() + defer t.mu.Unlock() + return t.upFiles +} + +// snapshot returns both legs at once, so a line rendered from it is consistent. +func (t *legTotals) snapshot() (dlFiles int, dlBytes int64, upFiles int, upBytes int64) { + t.mu.Lock() + defer t.mu.Unlock() + return t.dlFiles, t.dlBytes, t.upFiles, t.upBytes +}