fix: report why a download failed, retry it, and stop a failing source

core's downloader logs a failed transfer and returns nil, and nothing here
installed a logger, so logctx handed out a nop and the reason was destroyed.
The size check in finish was all that survived, which reported a failed
two-gigabyte fetch as "short download: got 0 bytes" — a symptom with no cause
an operator can act on. The logger it reaches for now feeds a core that keeps
error entries and stores them on the element they belong to, so finish reports
the reason ahead of the byte count.

A failed download is also retried inside the run, up to three passes over
whatever is still missing and spaced minutes apart, because the failure this
repairs is transient: a dead connection takes every transfer in flight with it
and all of them are fetchable again afterwards, while leaving them to the next
run costs a full re-walk of the chat and a re-index of the remote first. An
item is counted once however many attempts it took, and a pass that recovers
reports nothing.

The download leg gains the breaker the upload leg already had. A source that
refuses everything refuses the rest of the list in milliseconds, so the pass
gives up after five failures in a row and the run exits 3 rather than spending
the entire outstanding list finding that out.
This commit is contained in:
tiennm99 committed 2026-09-08 18:20:24 +07:00
1 parent 8794f11bcb
commit efb29c9c39
14 files changed
+591 -52

No files matched your search

+15 -1
View File
@@ -134,7 +134,7 @@ a chat.
| 0 | complete |
| 1 | ran, but files remain — run again |
| 2 | usage error |
| 3 | remote or Telegram failure, including a destination that stopped accepting uploads |
| 3 | remote or Telegram failure, including either leg refusing transfer after transfer |
| 4 | stalled: files remain, none of which can ever be fetched |
| 130 / 143 | interrupted (SIGINT / SIGTERM) |
@@ -187,6 +187,20 @@ verifier, which matched on name and non-zero size — on the archive this was
built for it found six objects that had been counted complete for months, one
of them 221 MiB standing in for a 2 GiB video. Re-running repairs them.
**Failures.** A failed download is retried inside the run: up to three passes
over whatever is still missing, spaced a minute and then two apart. The failure
this exists for is the transient one — a connection that dies takes every
transfer in flight with it, and all of them are fetchable again minutes later —
and leaving them to the next run costs a full re-walk of the chat and a re-index
of the remote before a byte can move. Each failure is reported with the reason
Telegram gave for it, not just the byte count that arrived.
Either leg gives up once five transfers in a row fail. A source or destination
that is refusing everything will refuse the rest of the list too, in
milliseconds, so the run stops and exits 3 instead of spending the whole
outstanding list finding that out. A run that trips and then recovers on a retry
reports nothing.
## Replacing the shell pipeline
Earlier versions of this repo were three bash scripts — `run.sh`,
+7 -5
View File
@@ -213,11 +213,13 @@ func syncCmd(ctx context.Context, args []string) error {
final.Write(os.Stdout)
switch {
case errors.Is(runErr, pipeline.ErrDestinationFailing):
// Exit 3, not 1. The destination refused upload after upload, and it
// will refuse them next pass too — a driver retrying on "incomplete"
// would walk 18k messages and re-download gigabytes into a remote that
// cannot take a byte, indefinitely.
case errors.Is(runErr, pipeline.ErrDestinationFailing),
errors.Is(runErr, pipeline.ErrSourceFailing):
// Exit 3, not 1. One half of the run refused transfer after transfer,
// and it will refuse them next pass too — a driver retrying on
// "incomplete" would walk the whole chat and re-index the whole remote
// to achieve nothing, indefinitely. The run already retried the failures
// itself before giving up, so this is not a transient verdict.
return runErr
case final.Stalled():
return fmt.Errorf("%w: %d file(s) remain, none of which can be fetched",
+2 -5
View File
@@ -6,7 +6,9 @@ require (
github.com/gotd/td v0.140.0
github.com/iyear/tdl/core v0.20.4
github.com/rclone/rclone v1.75.1
github.com/vbauerster/mpb/v8 v8.16.1
go.etcd.io/bbolt v1.5.0
go.uber.org/zap v1.28.0
golang.org/x/sync v0.22.0
)
@@ -38,7 +40,6 @@ require (
github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d // indirect
github.com/adrg/xdg v0.5.3 // indirect
github.com/anchore/go-lzo v0.1.1 // indirect
github.com/andybalholm/brotli v1.2.2 // indirect
github.com/andybalholm/cascadia v1.3.4 // indirect
github.com/apache/arrow-go/v18 v18.7.0 // indirect
github.com/appscode/go-querystring v0.0.0-20170504095604-0126cfb3f1dc // indirect
@@ -127,7 +128,6 @@ require (
github.com/gorilla/schema v1.4.1 // indirect
github.com/gotd/contrib v0.20.0 // indirect
github.com/gotd/ige v0.3.0 // indirect
github.com/gotd/log v0.1.0 // indirect
github.com/gotd/neo v0.1.5 // indirect
github.com/hashicorp/errwrap v1.1.0 // indirect
github.com/hashicorp/go-cleanhttp v0.5.2 // indirect
@@ -184,7 +184,6 @@ require (
github.com/putdotio/go-putio/putio v0.0.0-20200123120452-16d982cac2b8 // indirect
github.com/rclone/Proton-API-Bridge v1.0.5 // indirect
github.com/rclone/go-proton-api v1.0.4 // indirect
github.com/refraction-networking/utls v1.8.2 // indirect
github.com/relvacode/iso8601 v1.7.0 // indirect
github.com/rfjakob/eme v1.2.0 // indirect
github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06 // indirect
@@ -206,7 +205,6 @@ require (
github.com/ulikunitz/xz v0.5.15 // indirect
github.com/unknwon/goconfig v1.0.0 // indirect
github.com/vbauerster/cupwriter v0.0.4 // indirect
github.com/vbauerster/mpb/v8 v8.16.1 // indirect
github.com/wk8/go-ordered-map/v2 v2.1.8 // indirect
github.com/xanzy/ssh-agent v0.3.3 // indirect
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect
@@ -223,7 +221,6 @@ require (
go.opentelemetry.io/otel/trace v1.44.0 // indirect
go.uber.org/atomic v1.11.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.uber.org/zap v1.28.0 // indirect
go.yaml.in/yaml/v2 v2.4.4 // indirect
go.yaml.in/yaml/v3 v3.0.5 // indirect
golang.org/x/crypto v0.56.0 // indirect
-12
View File
@@ -294,16 +294,10 @@ github.com/gotd/contrib v0.20.0 h1:1Wc4+HMQiIKYQuGHVwVksIx152HFTP6B5n88dDe0ZYw=
github.com/gotd/contrib v0.20.0/go.mod h1:P6o8W4niqhDPHLA0U+SA/L7l3BQHYLULpeHfRSePn9o=
github.com/gotd/ige v0.3.0 h1:4f6LEHWsVDLBG0bT9wWG2/9TZb5aWm265G8ZlTXmRRU=
github.com/gotd/ige v0.3.0/go.mod h1:FE9bTaQtvfArizAcZuI4sS6gXaEUBmixdUufVHoCKac=
github.com/gotd/log v0.1.0 h1:4LJUEvafD1xtBwx2QkrlzFnRgbYXTlWqJPDi8BvrLbU=
github.com/gotd/log v0.1.0/go.mod h1:5ilhdu1Ux0QvDY/FF3Ojfw24Ws3SlCtyLwOpXy8KYXs=
github.com/gotd/log/logzap v0.1.1 h1:O6l7d8HUbODe+UMcrM47eXYDwdJ6RNmpQejLjrlcEIQ=
github.com/gotd/log/logzap v0.1.1/go.mod h1:5ObZkITbfhbsBOLzBkzmMk9QxXc0eNQpimau7zRL+Y8=
github.com/gotd/neo v0.1.5 h1:oj0iQfMbGClP8xI59x7fE/uHoTJD7NZH9oV1WNuPukQ=
github.com/gotd/neo v0.1.5/go.mod h1:9A2a4bn9zL6FADufBdt7tZt+WMhvZoc5gWXihOPoiBQ=
github.com/gotd/td v0.140.0 h1:trNBzTnhNtNwHsFp5qwKnNxQRAZJ6/BRE+uH3Lojauk=
github.com/gotd/td v0.140.0/go.mod h1:0ZkRxG7N+5ooG7/zdRXcnGautGPM6IKmyPQvdsAeF20=
github.com/gotd/td v0.161.0 h1:krbzsb70cakdrqF+MUIo+W7BkQTVhyB1kNS7X/+BLcY=
github.com/gotd/td v0.161.0/go.mod h1:7HdCs+zeJugdgZAF5iG8f70eOJvuiH2QzjoyUcysXbY=
github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
github.com/hashicorp/errwrap v1.1.0 h1:OxrOeh75EUXMY8TBjag2fzXGZ40LB6IKw45YeGUDY2I=
github.com/hashicorp/errwrap v1.1.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
@@ -386,8 +380,6 @@ github.com/mattn/go-colorable v0.1.15/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stg
github.com/mattn/go-isatty v0.0.23 h1:cYwCQTQf3HB6xUC+BtyCLZNr7IzbOmoZbmssVNzSyiQ=
github.com/mattn/go-isatty v0.0.23/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A=
github.com/mattn/go-runewidth v0.0.3/go.mod h1:LwmH8dsx7+W8Uxz3IHJYH5QSwggIsqBzpuz5H//U1FU=
github.com/mattn/go-runewidth v0.0.24 h1:cpokDiIn0MGnhdHwuWnJBITySJ20QyNGnY2kR/ay2DU=
github.com/mattn/go-runewidth v0.0.24/go.mod h1:XBkDxAl56ILZc9knddidhrOlY5R/pDhgLpndooCuJAs=
github.com/mattn/go-runewidth v0.0.28 h1:rPyg2ybwEKPebvpzVWe1gKBkH8EQFkxO4Y0hjBeLaBU=
github.com/mattn/go-runewidth v0.0.28/go.mod h1:3qAiGCV4Koz/yuveO58qUefmUTRm8r0IGEXZ9jeHp/8=
github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y=
@@ -461,8 +453,6 @@ github.com/rclone/go-proton-api v1.0.4 h1:AJW0e9pB4j0hVK4WqyGErFwaI+5MUQWPCtj5FY
github.com/rclone/go-proton-api v1.0.4/go.mod h1:QAlkFfswzrBuxvCORWV8rZdddg52hahMN98CFWoFW1E=
github.com/rclone/rclone v1.75.1 h1:kIxQcoDLj2Gke/gMSHK7OnxhX1Gu1cJBLP1kJZoaFp0=
github.com/rclone/rclone v1.75.1/go.mod h1:4zmMjGatCkSJPRZDpo+7y3xOl8S29EMUyKvZop5mHr4=
github.com/refraction-networking/utls v1.8.2 h1:j4Q1gJj0xngdeH+Ox/qND11aEfhpgoEvV+S9iJ2IdQo=
github.com/refraction-networking/utls v1.8.2/go.mod h1:jkSOEkLqn+S/jtpEHPOsVv/4V4EVnelwbMQl4vCWXAM=
github.com/relvacode/iso8601 v1.7.0 h1:BXy+V60stMP6cpswc+a93Mq3e65PfXCgDFfhvNNGrdo=
github.com/relvacode/iso8601 v1.7.0/go.mod h1:FlNp+jz+TXpyRqgmM7tnzHHzBnz776kmAH2h3sZCn0I=
github.com/rfjakob/eme v1.2.0 h1:8dAHL+WVAw06+7DkRKnRiFp1JL3QjcJEZFqDnndUaSI=
@@ -540,8 +530,6 @@ github.com/wk8/go-ordered-map/v2 v2.1.8 h1:5h/BUHu93oj4gIdvHHHGsScSTMijfx5PeYkE/
github.com/wk8/go-ordered-map/v2 v2.1.8/go.mod h1:5nJHM5DyteebpVlHnWMV0rPz6Zp7+xBAnxjb1X5vnTw=
github.com/xanzy/ssh-agent v0.3.3 h1:+/15pJfg/RsTxqYcX6fHqOXZwwMP+2VyYWJeWM2qQFM=
github.com/xanzy/ssh-agent v0.3.3/go.mod h1:6dzNDKs0J9rVPHPhaGCukekBHKqfl+L3KghI1Bc68Uw=
github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU=
github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E=
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 h1:ilQV1hzziu+LLM3zUTJ0trRztfwgjqKnBWNtSRkbmwM=
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78/go.mod h1:aL8wCCfTfSfmXjznFBSZNN13rSJjlIOI1fUNAtF7rmI=
github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74=
+80
View File
@@ -0,0 +1,80 @@
package pipeline
import (
"context"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
"github.com/iyear/tdl/core/logctx"
)
// captureCauses points core's downloader logging back into this package.
//
// core's Download never returns why a transfer failed. It logs the error and
// returns nil (downloader.go:47-60), so the size check in finish is the only
// evidence anything went wrong — and a size check can say how many bytes
// arrived, never why the rest did not. Nothing here installed a logger either,
// which means logctx handed out zap.NewNop() and the reason was destroyed: a
// failed two-gigabyte fetch was reported as "short download: got 0 bytes" with
// the flood wait, expired file reference, or dead connection behind it gone.
//
// So the error is taken from the one place it exists. The entry carries the
// element it belongs to, which is this package's own *elem, so the error is
// stored on it and finish reports it beside the size.
func captureCauses(ctx context.Context) context.Context {
return logctx.With(ctx, zap.New(&causeCore{}))
}
// causeCore is a zapcore.Core that keeps error entries and drops everything
// else. Fields are matched on their type rather than their key so a rename
// upstream cannot silently switch the capture off.
type causeCore struct{ fields []zapcore.Field }
func (c *causeCore) Enabled(l zapcore.Level) bool { return l >= zapcore.ErrorLevel }
func (c *causeCore) With(fields []zapcore.Field) zapcore.Core {
joined := make([]zapcore.Field, 0, len(c.fields)+len(fields))
joined = append(joined, c.fields...)
joined = append(joined, fields...)
return &causeCore{fields: joined}
}
func (c *causeCore) Check(e zapcore.Entry, ce *zapcore.CheckedEntry) *zapcore.CheckedEntry {
if c.Enabled(e.Level) {
return ce.AddCore(e, c)
}
return ce
}
// Write pairs a logged error with the element it was logged for.
//
// The write happens on the download worker's own goroutine, immediately before
// the deferred OnDone that reads it, so storing the error on the element needs
// no synchronisation.
func (c *causeCore) Write(_ zapcore.Entry, fields []zapcore.Field) error {
var (
el *elem
err error
)
for _, set := range [][]zapcore.Field{c.fields, fields} {
for _, f := range set {
switch f.Type {
case zapcore.ErrorType:
if e, ok := f.Interface.(error); ok {
err = e
}
case zapcore.ReflectType:
if e, ok := f.Interface.(*elem); ok {
el = e
}
}
}
}
if el != nil && err != nil {
el.cause = err
}
return nil
}
func (c *causeCore) Sync() error { return nil }
+39 -1
View File
@@ -38,6 +38,16 @@ type DownloadOptions struct {
// stop, when set, ends iteration cleanly from another goroutine — used to
// halt downloads once the destination has stopped accepting uploads.
stop *atomic.Bool
// maxFailures ends the pass after this many consecutive failures. Zero means
// no breaker. onTrip, when set, is called once if that happens — the caller
// decides what an abandoned pass means, because a retry that then succeeds
// must not report the source as down.
maxFailures int
onTrip func()
// seed starts the run's counters from an earlier pass's totals, so a retry
// continues the numbers a reporter is already displaying rather than
// restarting them at zero.
seed Stats
// 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)
@@ -55,7 +65,11 @@ type DownloadOptions struct {
// 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.
// and the rest continue, matching what a partial `tdl dl` pass did. Past
// maxFailures failures in a row the pass gives up on the remaining items and
// calls onTrip, because a source that is refusing every file will refuse the
// rest of the list too, in milliseconds, and a pass with no brake spends the
// whole todo list proving it.
func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
if o.Threads <= 0 {
o.Threads = 4
@@ -67,6 +81,11 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo
return nil, Stats{}, fmt.Errorf("create staging directory: %w", err)
}
// core's downloader logs why a transfer failed rather than returning it, so
// the logger it reaches for is pointed back into this package before any of
// its work starts. Without this the reason goes to a nop logger.
ctx = captureCauses(ctx)
it := newElemIter(seq, o.Staging, o.Takeout)
it.acquire = o.acquire
it.release = o.release
@@ -88,6 +107,12 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo
}
return ferr
}, o.Events)
// Set rather than passed: both are this package's own bookkeeping, and no
// worker exists yet, so the mutex the fields normally live under has
// nothing to protect them from here.
prog.stats = o.seed
prog.maxStreak = o.maxFailures
prog.onTrip = func() { it.stopped.Store(true) }
err := downloader.New(downloader.Options{
Pool: o.Pool,
@@ -103,6 +128,12 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo
if err == nil {
err = it.failure
}
// A tripped breaker is not an error here. It is reported to the caller, who
// knows whether another pass is coming; the items the pass never reached are
// simply absent from the outcomes.
if prog.brokeCircuit() && o.onTrip != nil {
o.onTrip()
}
// 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.
@@ -131,6 +162,13 @@ func finish(staging string, e *elem, downloadErr error) error {
if downloadErr == nil {
if err := checkSize(part, e.item.Size()); err != nil {
// The size is the symptom. e.cause is the reason, when the
// downloader logged one, and it is put first: "got 0 bytes" alone
// says a two-gigabyte fetch failed without saying anything an
// operator can act on.
if e.cause != nil {
err = fmt.Errorf("%w: %w", e.cause, err)
}
downloadErr = err
}
}
+104
View File
@@ -10,7 +10,9 @@ import (
"github.com/gotd/td/tg"
"github.com/iyear/tdl/core/downloader"
"github.com/iyear/tdl/core/logctx"
"github.com/iyear/tdl/core/tmedia"
"go.uber.org/zap"
"github.com/tiennm99dev/telegram-exporter/internal/naming"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
@@ -429,3 +431,105 @@ func TestElemIterReturnsReservationWhenOpenFails(t *testing.T) {
t.Errorf("acquired %d bytes but released %d — the reservation leaked", acquired, released)
}
}
// The reason a transfer failed only exists in core's log call, so the capture is
// tested through that exact call rather than by poking the field directly: it is
// the shape of the log entry — a reflected element and a zap error — that this
// depends on, and a change to it must fail here rather than silently go back to
// reporting a byte count with no reason.
func TestFinishReportsWhyTheTransferFailed(t *testing.T) {
staging := t.TempDir()
it := testItem(t, 5, "video.mp4", 1000)
e := openElem(t, staging, it)
want := errors.New("rpc error code 420: FLOOD_WAIT (60)")
ctx := captureCauses(t.Context())
logctx.From(ctx).Error("Download error",
zap.Any("element", downloader.Elem(e)), zap.Error(want))
// nil, exactly as core reports a transfer it has already logged and given up
// on, with nothing written to the file.
err := finish(staging, e, nil)
if !errors.Is(err, want) {
t.Fatalf("finish = %v, want the logged transfer error", err)
}
if !strings.Contains(err.Error(), "expected 1000") {
t.Errorf("the byte count must survive alongside the reason, got: %v", err)
}
}
// Below error level nothing is captured, so an ordinary debug entry cannot
// attach a bogus reason to an item that merely arrived short.
func TestCaptureIgnoresNonErrorLogging(t *testing.T) {
staging := t.TempDir()
it := testItem(t, 6, "video.mp4", 1000)
e := openElem(t, staging, it)
ctx := captureCauses(t.Context())
logctx.From(ctx).Debug("Start download elem", zap.Any("elem", downloader.Elem(e)))
if e.cause != nil {
t.Fatalf("cause = %v, want nil for a debug entry", e.cause)
}
if err := finish(staging, e, nil); !strings.Contains(err.Error(), "short download") {
t.Errorf("finish = %v, want the plain size failure", err)
}
}
// The breaker exists because a source that has stopped serving files fails every
// item in milliseconds: without it a pass spends the entire remaining todo list
// proving the same point. A success in between is what tells a run of bad luck
// apart from a source that is down, so it resets the streak.
func TestProgressBreakerTripsOnConsecutiveFailuresOnly(t *testing.T) {
staging := t.TempDir()
trips := 0
p := newProgress(func(*elem, error) error { return nil }, nil)
p.maxStreak = 3
p.onTrip = func() { trips++ }
fail := func(id int) { p.OnDone(openElem(t, staging, testItem(t, id, "a.mp4", 10)), errors.New("boom")) }
ok := func(id int) { p.OnDone(openElem(t, staging, testItem(t, id, "a.mp4", 10)), nil) }
fail(1)
fail(2)
ok(3) // the streak is broken here, so the next two must not trip it
fail(4)
fail(5)
if trips != 0 {
t.Fatalf("tripped after a success reset the streak (trips = %d)", trips)
}
if p.brokeCircuit() {
t.Fatal("brokeCircuit() = true before the limit was reached")
}
fail(6)
if trips != 1 {
t.Fatalf("trips = %d, want 1 at three failures in a row", trips)
}
if !p.brokeCircuit() {
t.Error("brokeCircuit() = false after the breaker fired")
}
// Fires once: the callback stops the iterator, and repeating it for every
// item still in flight would be noise.
fail(7)
if trips != 1 {
t.Errorf("trips = %d, want the breaker to fire exactly once", trips)
}
}
// With no limit set there is no breaker at all, which is what a standalone
// Download must keep doing: every item gets an attempt.
func TestProgressWithoutABreakerNeverTrips(t *testing.T) {
staging := t.TempDir()
p := newProgress(func(*elem, error) error { return nil }, nil)
p.onTrip = func() { t.Error("onTrip called with no limit configured") }
for id := 1; id <= 20; id++ {
p.OnDone(openElem(t, staging, testItem(t, id, "a.mp4", 10)), errors.New("boom"))
}
if p.brokeCircuit() {
t.Error("brokeCircuit() = true with no limit configured")
}
}
+7
View File
@@ -35,6 +35,13 @@ type elem struct {
item tgsource.Item
file *os.File
takeout bool
// cause is why the transfer failed, when the downloader logged a reason.
// It is the only route that reason has to reach finish, because core's
// Download logs it and returns nil; see captureCauses. Written by the log
// call on the download worker's goroutine and read by that worker's OnDone,
// so it needs no lock.
cause error
}
func (e *elem) File() downloader.File { return mediaFile{e.item} }
+12 -1
View File
@@ -1,6 +1,10 @@
package pipeline
import "github.com/tiennm99dev/telegram-exporter/internal/tgsource"
import (
"time"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// Events receives a run's per-item lifecycle.
//
@@ -24,6 +28,12 @@ type Events interface {
UploadStart(it tgsource.Item)
UploadDone(it tgsource.Item, err error)
// Retry announces another pass over the items the last one failed, and how
// long the run waits first. It is the one event that is not per-item, and it
// exists because the wait is measured in minutes: a display that went quiet
// for that long with no explanation is indistinguishable from a hang.
Retry(attempt, files int, wait time.Duration)
}
// nopEvents is used when a caller wants no reporting, so nothing on the hot
@@ -36,3 +46,4 @@ func (nopEvents) DownloadBytes(tgsource.Item, int64) {}
func (nopEvents) DownloadDone(tgsource.Item, error) {}
func (nopEvents) UploadStart(tgsource.Item) {}
func (nopEvents) UploadDone(tgsource.Item, error) {}
func (nopEvents) Retry(int, int, time.Duration) {}
+156 -17
View File
@@ -9,6 +9,7 @@ import (
"path/filepath"
"sync"
"sync/atomic"
"time"
"github.com/rclone/rclone/fs"
"golang.org/x/sync/semaphore"
@@ -72,6 +73,37 @@ func (r Result) Failed() []Outcome {
// gigabytes for nothing every time.
var ErrDestinationFailing = errors.New("destination stopped accepting uploads")
// ErrSourceFailing marks a run stopped because Telegram failed download after
// download. It is the download leg's counterpart to ErrDestinationFailing and
// means the same thing to a driver — another pass will fail the same way, so
// retrying on "incomplete" only re-walks the chat for nothing — but it points at
// the other half of the run, which is what an operator needs to know first.
//
// It survives the run's own retries: a pass that trips and then recovers reports
// nothing, so this only appears when the last attempt was still failing.
var ErrSourceFailing = errors.New("telegram stopped serving downloads")
// downloadAttempts is how many passes a run makes over the items it failed to
// fetch.
//
// Retrying inside the run is worth far more than leaving it to the next one: a
// fresh sync re-walks every message and re-indexes the whole remote before it
// can fetch a byte, while a retry here already has both. The failure this exists
// for is the transient one — a dead connection or a source that stops serving
// files for a few minutes takes down every transfer in flight, and every one of
// them is fetchable again afterwards.
const downloadAttempts = 3
// retryDelay spaces the passes out, and is a var so tests need not sleep.
//
// Minutes rather than seconds: the client's own recovery gives up on a broken
// connection only after the reconnect timeout (five minutes by default), so a
// pass that starts seconds after the last one failed is a pass into the same
// dead connection.
var retryDelay = func(attempt int) time.Duration {
return time.Duration(attempt) * time.Minute
}
// maxRecordedErrors bounds what a run keeps from a failing remote. Past this,
// the pattern is established and joining thousands of identical strings just
// makes the final message unreadable.
@@ -187,23 +219,97 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R
}()
}
dlOutcomes, stats, dlErr := download(ctx, seq, DownloadOptions{
Pool: o.Pool,
Staging: o.Staging,
Threads: o.Threads,
Limit: o.Limit,
Takeout: o.Takeout,
Events: o.Events,
acquire: budget.acquire,
release: budget.release,
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())
},
stop: stopDownloads,
})
var (
dlOutcomes []Outcome
stats Stats
dlErrs []error
at = make(map[int]int) // message id -> its place in dlOutcomes
)
pending := seq
for attempt := 1; ; attempt++ {
sourceDown := false
passOutcomes, passStats, err := download(ctx, pending, DownloadOptions{
Pool: o.Pool,
Staging: o.Staging,
Threads: o.Threads,
Limit: o.Limit,
Takeout: o.Takeout,
Events: o.Events,
acquire: budget.acquire,
release: budget.release,
maxFailures: o.MaxFailures,
seed: stats,
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())
},
stop: stopDownloads,
onTrip: func() { sourceDown = true },
})
stats = passStats
if err != nil {
dlErrs = append(dlErrs, err)
}
// One outcome per item, whatever it took: a later attempt replaces the
// earlier verdict in place, so an item that failed once and then
// arrived is reported as archived rather than as both.
for _, oc := range passOutcomes {
if i, ok := at[oc.Item.MessageID]; ok {
dlOutcomes[i] = oc
continue
}
at[oc.Item.MessageID] = len(dlOutcomes)
dlOutcomes = append(dlOutcomes, oc)
}
// Read, and the shared stop flag cleared, under the upload leg's own
// lock. The flag is what the download breaker sets to end a pass, so a
// retry needs it clear — but clearing it after the upload breaker has
// tripped in this same window would restart downloads into a
// destination that has stopped accepting them, and that leg never sets
// the flag twice.
mu.Lock()
destDown := tripped
if !destDown {
stopDownloads.Store(false)
}
mu.Unlock()
// Another pass is worth making only for items that failed on their own.
// A pass that returned an error failed as a whole — the walk broke, a
// name could not be opened, the run was cancelled — and the items it
// never reached are lost to this run either way, so retrying the few
// that failed first would dress that up as a nearly complete run. A
// destination that has stopped accepting uploads rules it out too:
// fetching more would only fill staging with files it will refuse.
again := failedItems(passOutcomes)
if len(again) > 0 && attempt < downloadAttempts &&
err == nil && !destDown && ctx.Err() == nil {
// These items are being fetched again, so their first attempt comes
// back out of the totals: one item is one file to fetch, not one per
// attempt. Bytes already transferred stay counted — they were really
// spent, and throughput is the honest figure.
stats = rollback(stats, again)
wait := retryDelay(attempt)
o.Events.Retry(attempt+1, len(again), wait)
if serr := sleep(ctx, wait); serr == nil {
pending = itemsSeq(again)
continue
}
dlErrs = append(dlErrs, fmt.Errorf("cancelled before retrying %d download(s): %w",
len(again), ctx.Err()))
}
if sourceDown {
dlErrs = append(dlErrs, fmt.Errorf("%w: gave up after %d consecutive download failure(s)",
ErrSourceFailing, o.MaxFailures))
}
break
}
dlErr := errors.Join(dlErrs...)
// Safe only because Download joined its workers, which is guaranteed by
// elemIter.Err always being nil.
@@ -232,6 +338,39 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R
return Result{Stats: stats, Outcomes: dlOutcomes}, errors.Join(dlErr, upErr)
}
// failedItems lists the items a pass did not manage to fetch.
func failedItems(outcomes []Outcome) []tgsource.Item {
var out []tgsource.Item
for _, oc := range outcomes {
if oc.Err != nil {
out = append(out, oc.Item)
}
}
return out
}
// rollback takes a pass's failures back out of the running totals so the next
// pass counts them once, not twice.
func rollback(s Stats, again []tgsource.Item) Stats {
for _, it := range again {
s.Started--
s.Failed--
s.BytesTotal -= it.Size()
}
return s
}
// itemsSeq feeds a retry pass the items the last one failed.
func itemsSeq(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
}
}
}
}
// 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 }
+31
View File
@@ -35,6 +35,17 @@ type progress struct {
outcomes []Outcome
inFlight map[int]int64 // message id -> bytes written so far
// maxStreak is how many failures in a row end the pass; zero disables the
// breaker. streak counts them, and onTrip is called once when the limit is
// reached. The download leg needs this for the same reason the upload leg
// does: when the source stops serving files every item fails in
// milliseconds, and without a breaker a pass burns the entire remaining
// todo list on transfers that cannot succeed.
maxStreak int
streak int
tripped bool
onTrip func()
finish func(*elem, error) error
events Events
}
@@ -88,16 +99,36 @@ func (p *progress) OnDone(e downloader.Elem, err error) {
delete(p.inFlight, el.item.MessageID)
if err != nil {
p.stats.Failed++
p.streak++
} else {
p.stats.Done++
p.streak = 0
}
p.outcomes = append(p.outcomes, Outcome{Item: el.item, Err: err})
stats := p.stats
// Decided under the lock and acted on outside it: onTrip stops the
// iterator, and holding this mutex across it would put the callback's
// locking order inside this one's.
trip := p.maxStreak > 0 && p.streak >= p.maxStreak && !p.tripped
if trip {
p.tripped = true
}
p.mu.Unlock()
if trip && p.onTrip != nil {
p.onTrip()
}
p.events.DownloadDone(el.item, err)
p.events.Stats(stats)
}
// brokeCircuit reports whether the failure streak ended this pass.
func (p *progress) brokeCircuit() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.tripped
}
func (p *progress) results() ([]Outcome, Stats) {
p.mu.Lock()
defer p.mu.Unlock()
+124 -10
View File
@@ -9,6 +9,7 @@ import (
"path/filepath"
"sync/atomic"
"testing"
"time"
_ "github.com/rclone/rclone/backend/local"
"github.com/rclone/rclone/fs"
@@ -25,11 +26,18 @@ import (
// pieces are wired together, and a regression in any of them is silent.
// runFake substitutes the download step for the duration of a test.
//
// It also collapses the wait between retry passes. A run retries the items it
// failed to fetch, and at the real cadence every test with a failure in it would
// spend minutes asleep.
func runFake(t *testing.T, f func(context.Context, iter.Seq2[tgsource.Item, error], DownloadOptions) ([]Outcome, Stats, error)) {
t.Helper()
prev := download
prev, prevDelay := download, retryDelay
download = f
t.Cleanup(func() { download = prev })
retryDelay = func(int) time.Duration { return time.Millisecond }
t.Cleanup(func() {
download, retryDelay = prev, prevDelay
})
}
// stageItems is a download step that writes each item's bytes into staging and
@@ -37,10 +45,10 @@ func runFake(t *testing.T, f func(context.Context, iter.Seq2[tgsource.Item, erro
// 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
)
var outcomes []Outcome
// Seeded exactly as the real step does, so a retry pass continues the
// run's totals instead of restarting them.
stats := o.seed
for it, err := range seq {
if err != nil {
return outcomes, stats, err
@@ -213,10 +221,10 @@ func TestRunStopsDownloadingAfterConsecutiveUploadFailures(t *testing.T) {
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
)
var outcomes []Outcome
// Seeded exactly as the real step does, so a retry pass continues the
// run's totals instead of restarting them.
stats := o.seed
for it := range seqValues(seq) {
if o.stop.Load() {
break
@@ -305,3 +313,109 @@ func assertEmpty(t *testing.T, dir string) {
t.Errorf("staging is not empty: %v", names)
}
}
// A failure that clears on a later pass has to be reported as an archived file,
// not as both a failure and a success. The whole point of retrying in-run is
// that the outage which cost these items is usually over minutes later, and the
// next sync would have to re-walk the chat and re-index the remote to find out.
func TestRunRetriesFailedDownloadsWithinTheRun(t *testing.T) {
const size = 512
items := testItems(3, size)
var passes atomic.Int64
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
// The middle item fails on the first pass only, as a transient outage
// looks from here.
fail := map[int]bool{2: passes.Add(1) == 1}
return stageItems(fail)(ctx, seq, o)
})
o, staging, dstDir := runOpts(t, 0)
res, err := Run(t.Context(), seqOf(items), o)
if err != nil {
t.Fatalf("Run: %v", err)
}
if passes.Load() != 2 {
t.Errorf("passes = %d, want 2 — the failure was not retried", passes.Load())
}
if got := len(res.Failed()); got != 0 {
t.Errorf("Failed() = %d, want 0: %v", got, res.Failed())
}
if len(res.Outcomes) != len(items) {
t.Errorf("Outcomes = %d, want %d — a retried item must not be reported twice",
len(res.Outcomes), len(items))
}
// Counters follow the same rule: an item is one file to fetch however many
// attempts it took, or the closing summary claims more work than the chat
// contains.
if res.Stats.Done != 3 || res.Stats.Failed != 0 || res.Stats.Started != 3 {
t.Errorf("Stats = %+v, want 3 started, 3 done, 0 failed", res.Stats)
}
for _, it := range items {
if _, serr := os.Stat(filepath.Join(dstDir, it.Name)); serr != nil {
t.Errorf("%s should be on the destination: %v", it.Name, serr)
}
}
assertEmpty(t, staging)
}
// Retrying is bounded. A source that is down stays down for the run, and the
// verdict has to be the distinct sentinel: a driver looping on "incomplete"
// would otherwise re-walk the whole chat forever against a dead source.
func TestRunGivesUpAndNamesTheSourceAfterEveryPassFails(t *testing.T) {
const size = 512
items := testItems(2, size)
var passes atomic.Int64
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
passes.Add(1)
outcomes, stats, err := stageItems(map[int]bool{1: true, 2: true})(ctx, seq, o)
// What Download reports once its own breaker has ended the pass.
if o.onTrip != nil {
o.onTrip()
}
return outcomes, stats, err
})
o, staging, _ := runOpts(t, 0)
o.MaxFailures = 2
res, err := Run(t.Context(), seqOf(items), o)
if !errors.Is(err, ErrSourceFailing) {
t.Fatalf("Run = %v, want ErrSourceFailing", err)
}
if got := passes.Load(); got != downloadAttempts {
t.Errorf("passes = %d, want %d", got, downloadAttempts)
}
if got := len(res.Failed()); got != len(items) {
t.Errorf("Failed() = %d, want %d", got, len(items))
}
assertEmpty(t, staging)
}
// A trip that the retry then clears must not reach the caller. Reported anyway,
// it would send a driver to the "source is down, stop" exit code on a run that
// finished everything.
func TestRunDoesNotReportASourceThatRecovered(t *testing.T) {
const size = 512
items := testItems(2, size)
var passes atomic.Int64
runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) {
first := passes.Add(1) == 1
outcomes, stats, err := stageItems(map[int]bool{1: first, 2: first})(ctx, seq, o)
if first && o.onTrip != nil {
o.onTrip()
}
return outcomes, stats, err
})
o, _, _ := runOpts(t, 0)
o.MaxFailures = 2
res, err := Run(t.Context(), seqOf(items), o)
if err != nil {
t.Fatalf("Run = %v, want nil after the retry succeeded", err)
}
if got := len(res.Failed()); got != 0 {
t.Errorf("Failed() = %d, want 0", got)
}
}
+7
View File
@@ -210,6 +210,13 @@ func (l *Live) UploadDone(it tgsource.Item, err error) {
l.upTotal.SetCurrent(l.legs.addUpload(it.Size()))
}
// Retry prints above the bars rather than through them, so the pass that is
// about to start is announced without shredding the display.
func (l *Live) Retry(attempt, files int, wait time.Duration) {
fmt.Fprintf(l.p, " retrying %s file(s) in %s — attempt %d\n",
humanCount(files), wait.Round(time.Second), attempt)
}
// Finish drains the bars and prints the closing summary.
func (l *Live) Finish(s pipeline.Stats) {
l.mu.Lock()
+7
View File
@@ -111,6 +111,13 @@ func (r *Reporter) UploadDone(it tgsource.Item, err error) {
fmt.Fprintf(r.w, " archived %-10s %q\n", humanBytes(it.Size()), it.Name)
}
func (r *Reporter) Retry(attempt, files int, wait time.Duration) {
r.mu.Lock()
defer r.mu.Unlock()
fmt.Fprintf(r.w, " retrying %s file(s) in %s — attempt %d\n",
humanCount(files), wait.Round(time.Second), attempt)
}
// Finish writes the closing summary.
func (r *Reporter) Finish(s pipeline.Stats) {
r.mu.Lock()