diff --git a/README.md b/README.md index a9d062b..07382ff 100644 --- a/README.md +++ b/README.md @@ -187,6 +187,28 @@ 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. +**File references.** Telegram hands out a short-lived token with every file +location, and a run over a large chat outlives the ones its walk collected — the +lifetime is undocumented, but an observed run stopped just under two hours in +with `FILE_REFERENCE_EXPIRED` on every remaining file. So each message is re-read +for a live token immediately before its own download, which costs one round trip +per transfer and needs no guess at how long a token lasts. Retry passes go +through the same path, so a token that dies during a single very large transfer +is replaced rather than replayed. + +The re-read has to describe the same file — same name, same length — or it is not +the file this run recorded, and archiving it under that name would store the +wrong bytes. A message that fails that check, or that has been deleted, or that +no longer holds a file at all, is skipped and reported: nothing can be archived +for it, one of them must not strand the rest of the chat, and the closing +`verify` still lists it as outstanding. + +A re-read that fails for any other reason is counted as a failed download, not a +skip, because that is what it is — the connection that serves the re-read is the +one that serves the file. So it is retried with the rest, it shows up in the +progress report and the closing failure list, and enough of them in a row trip +the same breaker below. + **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 diff --git a/cmd/tgexport/sync.go b/cmd/tgexport/sync.go index 1b9f306..5de1301 100644 --- a/cmd/tgexport/sync.go +++ b/cmd/tgexport/sync.go @@ -176,6 +176,10 @@ func syncCmd(ctx context.Context, args []string) error { Confirm: *confirm, Takeout: *takeout, Events: rep, + // The walk above collected every item's file reference up front, and + // a run this size spends longer downloading than one of those stays + // valid, so each is re-read immediately before its own download. + Refresh: tgsource.Refresher(api, peer), // Re-checked during the run, not only before it: an archive of this // size runs for hours, and the destination can fill in the middle. FreeBytes: func(ctx context.Context) (int64, bool) { diff --git a/internal/pipeline/download.go b/internal/pipeline/download.go index 8ae7130..5fba176 100644 --- a/internal/pipeline/download.go +++ b/internal/pipeline/download.go @@ -29,6 +29,11 @@ type DownloadOptions struct { // concurrently. Events Events + // Refresh, when set, re-mints an item's file reference just before its + // download starts. Unset means the reference the walk minted is used as-is, + // which only holds up for a run shorter than a reference's lifetime. + Refresh tgsource.Refresh + // acquire reserves staging space before a download starts, blocking until // there is room. Unset means no bound. release hands a reservation back for // an item that never reaches a download. @@ -87,6 +92,7 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo ctx = captureCauses(ctx) it := newElemIter(seq, o.Staging, o.Takeout) + it.refresh = o.Refresh it.acquire = o.acquire it.release = o.release if o.stop != nil { @@ -113,6 +119,9 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo prog.stats = o.seed prog.maxStreak = o.maxFailures prog.onTrip = func() { it.stopped.Store(true) } + // Wired after prog exists, so an item the iterator cannot even hand over is + // counted, reported and retried like any other failed transfer. + it.refused = prog.refused err := downloader.New(downloader.Options{ Pool: o.Pool, diff --git a/internal/pipeline/download_test.go b/internal/pipeline/download_test.go index bd7372f..aa52767 100644 --- a/internal/pipeline/download_test.go +++ b/internal/pipeline/download_test.go @@ -3,6 +3,7 @@ package pipeline import ( "context" "errors" + "fmt" "os" "path/filepath" "strings" @@ -533,3 +534,260 @@ func TestProgressWithoutABreakerNeverTrips(t *testing.T) { t.Error("brokeCircuit() = true with no limit configured") } } + +// The reference the walk minted is dead by the time a long run reaches the end +// of its list, so the location the downloader is handed must be the one refresh +// just produced. Without that substitution every fetch past the reference's +// lifetime returns FILE_REFERENCE_EXPIRED and the breaker abandons the rest. +func TestElemIterDownloadsTheRefreshedLocation(t *testing.T) { + staging := t.TempDir() + stale := testItem(t, 42, "clip.mp4", 10) + + var asked []int + it := newElemIter(func(yield func(tgsource.Item, error) bool) { yield(stale, nil) }, + staging, false) + it.refresh = func(_ context.Context, in tgsource.Item) (tgsource.Item, error) { + asked = append(asked, in.MessageID) + fresh := *in.Media + fresh.InputFileLoc = &tg.InputDocumentFileLocation{ + ID: in.Media.InputFileLoc.(*tg.InputDocumentFileLocation).ID, + FileReference: []byte("fresh"), + } + in.Media = &fresh + return in, nil + } + defer func() { _ = it.Close() }() + + if !it.Next(t.Context()) { + t.Fatalf("Next produced nothing: %v", it.failure) + } + if len(asked) != 1 || asked[0] != 42 { + t.Fatalf("refresh calls = %v, want one for message 42", asked) + } + + loc, ok := it.Value().File().Location().(*tg.InputDocumentFileLocation) + if !ok { + t.Fatalf("location is %T, want *tg.InputDocumentFileLocation", it.Value().File().Location()) + } + if string(loc.FileReference) != "fresh" { + t.Errorf("downloader was handed file reference %q, want the refreshed one", loc.FileReference) + } +} + +// A message that cannot be re-read — deleted, or now holding a different file — +// must not strand every message behind it, which is what ending the walk would +// mean. It is reported instead, so nothing vanishes silently. +func TestElemIterSkipsItemsRefreshCannotFindAndKeepsGoing(t *testing.T) { + staging := t.TempDir() + gone := testItem(t, 7, "deleted.mp4", 10) + good := testItem(t, 8, "fine.mp4", 10) + + seq := func(yield func(tgsource.Item, error) bool) { + if !yield(gone, nil) { + return + } + yield(good, nil) + } + it := newElemIter(seq, staging, false) + it.refresh = func(_ context.Context, in tgsource.Item) (tgsource.Item, error) { + if in.MessageID == 7 { + return tgsource.Item{}, fmt.Errorf("%w: message 7 was deleted", tgsource.ErrGone) + } + return in, nil + } + defer func() { _ = it.Close() }() + + if !it.Next(t.Context()) { + t.Fatal("an unrefreshable item ended the walk; the item after it was never reached") + } + if got := it.current.item.MessageID; got != 8 { + t.Fatalf("Next yielded message %d, want the item after the missing one", got) + } + if it.failure != nil { + t.Errorf("a skipped item must not fail the run, got: %v", it.failure) + } + if len(it.skipped) != 1 { + t.Fatalf("skipped = %d, want 1", len(it.skipped)) + } + if !strings.Contains(it.skipped[0].Error(), "message 7") { + t.Errorf("the skip should name the message, got: %v", it.skipped[0]) + } +} + +// A cancelled run ends the walk rather than skipping its way through the rest of +// the list, one refused refresh at a time. +func TestElemIterStopsWhenRefreshIsCancelled(t *testing.T) { + staging := t.TempDir() + seq := seqOf(testItems(3, 10)) + + ctx, cancel := context.WithCancel(t.Context()) + it := newElemIter(seq, staging, false) + it.refused = func(tgsource.Item, error) { + t.Error("cancellation was recorded as an item failure, which the breaker would count") + } + it.refresh = func(rctx context.Context, in tgsource.Item) (tgsource.Item, error) { + cancel() + return tgsource.Item{}, fmt.Errorf("re-read message %d: %w", in.MessageID, rctx.Err()) + } + defer func() { _ = it.Close() }() + + if it.Next(ctx) { + t.Fatal("Next produced an item after cancellation") + } + if !errors.Is(it.failure, context.Canceled) { + t.Errorf("failure = %v, want context.Canceled", it.failure) + } + if len(it.skipped) != 0 { + t.Errorf("cancellation was recorded as %d skipped item(s), not a failure", len(it.skipped)) + } +} + +// A refresh that failed for any reason other than the message being gone is the +// same event as a transfer dying — Telegram is not serving this file — and has +// to be counted as one. An item that never reaches a worker never reaches OnAdd +// or OnDone, so on the skip path a source outage would move no counter, trip no +// breaker and produce no Outcome for Run to retry, and would walk the whole todo +// list one dead round trip at a time. +func TestElemIterRecordsRefreshFailuresAsFailuresNotSkips(t *testing.T) { + staging := t.TempDir() + items := []tgsource.Item{ + testItem(t, 1, "a.mp4", 10), + testItem(t, 2, "b.mp4", 20), + } + + it := newElemIter(seqOf(items), staging, false) + defer func() { _ = it.Close() }() + + prog := newProgress(func(*elem, error) error { return nil }, nil) + prog.maxStreak = 2 + var tripped bool + prog.onTrip = func() { tripped = true; it.stopped.Store(true) } + it.refused = prog.refused + it.refresh = func(context.Context, tgsource.Item) (tgsource.Item, error) { + return tgsource.Item{}, errors.New("dial telegram: connection refused") + } + + for it.Next(t.Context()) { + t.Fatal("Next handed over an item whose reference could not be refreshed") + } + + if it.failure != nil { + t.Errorf("failure = %v, want none — a refused item ends the pass, not the program", it.failure) + } + if len(it.skipped) != 0 { + t.Errorf("skipped = %d, want 0 — a transient re-read failure is not a skip: %v", + len(it.skipped), it.skipped) + } + + outcomes, stats := prog.results() + if len(outcomes) != 2 { + t.Fatalf("outcomes = %d, want 2 — Run builds its retry list from these", len(outcomes)) + } + for _, oc := range outcomes { + if oc.Err == nil { + t.Errorf("message %d recorded without an error", oc.Item.MessageID) + } + } + if got := failedItems(outcomes); len(got) != 2 { + t.Errorf("failedItems = %d, want 2 — these have to reach the next pass", len(got)) + } + if stats.Started != 2 || stats.Failed != 2 || stats.Done != 0 { + t.Errorf("Stats = %+v, want 2 started, 2 failed, 0 done", stats) + } + if stats.BytesTotal != 30 { + t.Errorf("BytesTotal = %d, want 30 — the report's target must include them", stats.BytesTotal) + } + if !tripped { + t.Error("two consecutive refused items did not trip the breaker") + } +} + +// A message that is gone is the one refresh failure that is genuinely permanent, +// so it stays on the skip path and must not spend the breaker's streak. +func TestElemIterDoesNotCountAMissingMessageAsAFailure(t *testing.T) { + staging := t.TempDir() + it := newElemIter(seqOf([]tgsource.Item{testItem(t, 1, "a.mp4", 10)}), staging, false) + defer func() { _ = it.Close() }() + + prog := newProgress(func(*elem, error) error { return nil }, nil) + prog.maxStreak = 1 + prog.onTrip = func() { t.Error("a gone message tripped the source breaker") } + it.refused = func(tgsource.Item, error) { t.Error("a gone message was recorded as a failure") } + it.refresh = func(_ context.Context, in tgsource.Item) (tgsource.Item, error) { + return tgsource.Item{}, fmt.Errorf("%w: message %d was deleted", tgsource.ErrGone, in.MessageID) + } + + for it.Next(t.Context()) { + } + if len(it.skipped) != 1 { + t.Fatalf("skipped = %d, want 1", len(it.skipped)) + } + if outcomes, _ := prog.results(); len(outcomes) != 0 { + t.Errorf("outcomes = %d, want 0 — a gone message is not worth retrying", len(outcomes)) + } +} + +// The refresh has to come after the reservation, not before it. That acquire is +// the run's brake: it blocks for as long as the destination is slow, and a +// reference minted before it would spend that whole wait ageing — which on a +// stalled remote is long enough to expire it again, reproducing the failure the +// refresh exists to prevent. +func TestElemIterRefreshesAfterReservingSpace(t *testing.T) { + staging := t.TempDir() + it := newElemIter(seqOf([]tgsource.Item{testItem(t, 1, "a.mp4", 4096)}), staging, false) + defer func() { _ = it.Close() }() + + var acquired int64 + it.acquire = func(_ context.Context, n int64) error { acquired += n; return nil } + it.release = func(n int64) { acquired -= n } + + var heldAtRefresh int64 + it.refresh = func(_ context.Context, in tgsource.Item) (tgsource.Item, error) { + heldAtRefresh = acquired + return in, nil + } + + if !it.Next(t.Context()) { + t.Fatalf("Next produced nothing: %v", it.failure) + } + if heldAtRefresh != 4096 { + t.Errorf("staging held %d bytes when the reference was refreshed, want 4096 — "+ + "the refresh ran before the acquire it is supposed to follow", heldAtRefresh) + } +} + +// Whichever way a refresh fails, the bytes it reserved go back. Without that the +// budget shrinks by that much for the rest of the run. +func TestElemIterReturnsReservationWhenRefreshFails(t *testing.T) { + staging := t.TempDir() + + for _, tc := range []struct { + name string + err error + }{ + {"gone", fmt.Errorf("%w: deleted", tgsource.ErrGone)}, + {"transient", errors.New("dial telegram: connection refused")}, + } { + t.Run(tc.name, func(t *testing.T) { + it := newElemIter(seqOf([]tgsource.Item{testItem(t, 1, "a.mp4", 4096)}), staging, false) + defer func() { _ = it.Close() }() + + var acquired, released int64 + it.acquire = func(_ context.Context, n int64) error { acquired += n; return nil } + it.release = func(n int64) { released += n } + it.refresh = func(context.Context, tgsource.Item) (tgsource.Item, error) { + return tgsource.Item{}, tc.err + } + + for it.Next(t.Context()) { + } + if acquired == 0 { + t.Fatal("nothing was reserved, so the test proves nothing") + } + if acquired != released { + t.Errorf("acquired %d bytes but released %d — the reservation leaked", + acquired, released) + } + }) + } +} diff --git a/internal/pipeline/elem.go b/internal/pipeline/elem.go index c56b3d6..27f5ea1 100644 --- a/internal/pipeline/elem.go +++ b/internal/pipeline/elem.go @@ -3,6 +3,7 @@ package pipeline import ( "context" + "errors" "fmt" "io" "iter" @@ -76,6 +77,14 @@ type elemIter struct { staging string takeout bool + // refresh re-mints an item's file reference just before its download; see + // reminted. Unset leaves the walk's own reference in place. + refresh tgsource.Refresh + // refused records an item that failed before any transfer could start, so it + // counts as a failure rather than disappearing. Unset means such an item is + // only skipped. + refused func(tgsource.Item, error) + // acquire reserves staging space for the next item. Blocking here is what // makes backpressure work: core's Download calls Next from its dispatch // loop (downloader.go:38), so a blocked Next stops new downloads starting @@ -122,7 +131,6 @@ func newElemIter(seq iter.Seq2[tgsource.Item, error], staging string, takeout bo } func (i *elemIter) Next(ctx context.Context) bool { - var item tgsource.Item for { if i.failure != nil || i.stopped.Load() { return false @@ -132,11 +140,7 @@ func (i *elemIter) Next(ctx context.Context) bool { return false } - var ( - err error - ok bool - ) - item, err, ok = i.next() + item, err, ok := i.next() if !ok { return false } @@ -151,35 +155,110 @@ func (i *elemIter) Next(ctx context.Context) bool { // is either a second caller or a gap there; in both cases one unwritable // name must not strand the rest of the chat behind it. if err := naming.Safe(item.Name); err != nil { - if len(i.skipped) < maxRecordedErrors { - i.skipped = append(i.skipped, fmt.Errorf("message %d: %w", item.MessageID, err)) + i.skip(item, err) + continue + } + + if i.acquire != nil { + if err := i.acquire(ctx, item.Size()); err != nil { + i.failure = err + return false + } + } + + // The refreshed item replaces the original only on success. A failed + // refresh reports a zero Item, and handing that to the reservation or to + // the failure record would release nothing and blame message 0. + fresh, err := i.reminted(ctx, item) + if err == nil { + item = fresh + } else { + // Nothing was staged and nothing will be, so the reservation goes + // back before anything else; otherwise the cap ends the run + // over-committed by every item that failed here. + i.giveBack(item) + if ctx.Err() != nil { + i.failure = err + return false + } + // A message that is gone is skipped, for the same reason an + // unwritable name is: it is a permanent, per-message fact, and one + // of them must not strand every message behind it. Nothing is + // archived for it, so the run's closing verify still reports it as + // outstanding. + // + // Anything else is a failure and is recorded as one. Routing it + // there rather than onto the skip path is what keeps the breaker, + // the run's retry passes and the progress report working: a re-read + // fails because the connection did, which is exactly when every + // download is failing too, and a skip is invisible to all three — a + // source outage would otherwise walk the whole todo list one dead + // round trip at a time and still exit as merely incomplete. + if errors.Is(err, tgsource.ErrGone) { + i.skip(item, err) + } else if i.refused != nil { + i.refused(item, err) } continue } - break - } - if i.acquire != nil { - if err := i.acquire(ctx, item.Size()); err != nil { - i.failure = err + f, err := os.OpenFile(partPath(i.staging, item), os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + // The reservation is handed back here because this item will never + // reach a download, so no OnDone will ever release it for us. + i.giveBack(item) + i.failure = fmt.Errorf("open destination for message %d: %w", item.MessageID, err) return false } - } + i.opened = append(i.opened, f) - f, err := os.OpenFile(partPath(i.staging, item), os.O_CREATE|os.O_RDWR, 0o600) - if err != nil { - // The reservation is handed back here because this item will never - // reach a download, so no OnDone will ever release it for us. - if i.release != nil { - i.release(item.Size()) - } - i.failure = fmt.Errorf("open destination for message %d: %w", item.MessageID, err) - return false + i.current = &elem{item: item, file: f, takeout: i.takeout} + return true } - i.opened = append(i.opened, f) +} - i.current = &elem{item: item, file: f, takeout: i.takeout} - return true +// reminted replaces an item's file reference with one Telegram has just issued. +// +// The reference in Media.InputFileLoc was minted when the walk saw the message, +// and Telegram expires those. The lifetime is undocumented and was observed to +// outlast an hour but not two, which is less than a run over a large chat spends +// downloading — so by the time the dispatch loop reaches an item near the end of +// the list its reference is dead, every remaining fetch returns +// FILE_REFERENCE_EXPIRED, and the breaker reads that as a dead source and +// abandons everything still to come. +// +// It is called as late as Next can leave it, and in particular after the budget +// acquire rather than before. core's downloader reads Elem.File().Location() +// once per attempt, inside the worker (downloader.go:89), so the only wait left +// between here and there is for a free worker slot, which one download bounds. +// Re-minting before the acquire would reintroduce the very failure this exists +// to prevent: that acquire is the run's brake and blocks for as long as the +// destination is slow, which on a stalled remote is long enough to expire a +// fresh token all over again. +// +// One round trip per transfer, and no guess at how long a reference lives. +// Retries get it for free: Run feeds a failed item straight back through +// Download, so a reference that died mid-transfer — a multi-gigabyte file can +// outlive its own token — is re-minted on the next pass rather than replayed. +func (i *elemIter) reminted(ctx context.Context, item tgsource.Item) (tgsource.Item, error) { + if i.refresh == nil { + return item, nil + } + return i.refresh(ctx, item) +} + +// skip records an item no run will ever fetch, without ending this one. +func (i *elemIter) skip(item tgsource.Item, err error) { + if len(i.skipped) < maxRecordedErrors { + i.skipped = append(i.skipped, fmt.Errorf("message %d: %w", item.MessageID, err)) + } +} + +// giveBack returns an item's staging reservation. +func (i *elemIter) giveBack(item tgsource.Item) { + if i.release != nil { + i.release(item.Size()) + } } func (i *elemIter) Value() downloader.Elem { return i.current } diff --git a/internal/pipeline/pipeline.go b/internal/pipeline/pipeline.go index 351b37c..df25140 100644 --- a/internal/pipeline/pipeline.go +++ b/internal/pipeline/pipeline.go @@ -46,6 +46,12 @@ type Options struct { // Events, when set, receives the run's per-item lifecycle: which files are // downloading, which are uploading, and how far along each one is. Events Events + + // Refresh, when set, re-reads a message for a live file reference just + // before its download starts. A run over a large chat outlasts the + // references its walk collected, so without this every fetch past that point + // fails with FILE_REFERENCE_EXPIRED; see elemIter.Next. + Refresh tgsource.Refresh } // Result is what a run achieved. @@ -235,6 +241,7 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R Limit: o.Limit, Takeout: o.Takeout, Events: o.Events, + Refresh: o.Refresh, acquire: budget.acquire, release: budget.release, maxFailures: o.MaxFailures, diff --git a/internal/pipeline/progress.go b/internal/pipeline/progress.go index dade2e2..ba354fe 100644 --- a/internal/pipeline/progress.go +++ b/internal/pipeline/progress.go @@ -122,6 +122,46 @@ func (p *progress) OnDone(e downloader.Elem, err error) { p.events.Stats(stats) } +// refused records an item that failed before any transfer could start. +// +// It exists so that "the message could not be re-read" counts as the same kind +// of event as "the transfer died", because operationally it is: both mean +// Telegram is not serving this file right now, and both are usually the same +// outage. An item that never reaches a worker never reaches OnAdd or OnDone, so +// without this it would move no counter, feed no streak, produce no Outcome for +// Run to retry, and leave the live report on a still frame while the pass walked +// the rest of the list. +// +// Everything OnDone does for a failure it does here, minus the file: the item is +// counted as started and failed, its bytes join the total the report is working +// towards, the streak advances so a dead source still trips the breaker, and the +// Outcome is what makes Run try the item again on the next pass. +func (p *progress) refused(it tgsource.Item, err error) { + p.mu.Lock() + p.stats.Started++ + p.stats.BytesTotal += it.Size() + p.stats.Failed++ + p.streak++ + p.outcomes = append(p.outcomes, Outcome{Item: it, Err: err}) + stats := p.stats + // Decided under the lock and acted on outside it, as in OnDone. + 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() + } + // Reported as a start and an immediate end rather than as an end alone: a + // reporter tracking which files are in flight is entitled to see every item + // it is told about finish, and one it never saw start would be a stray. + p.events.DownloadStart(it) + p.events.DownloadDone(it, err) + p.events.Stats(stats) +} + // brokeCircuit reports whether the failure streak ended this pass. func (p *progress) brokeCircuit() bool { p.mu.Lock() diff --git a/internal/pipeline/run_test.go b/internal/pipeline/run_test.go index d5b6c30..ce8ed1d 100644 --- a/internal/pipeline/run_test.go +++ b/internal/pipeline/run_test.go @@ -419,3 +419,35 @@ func TestRunDoesNotReportASourceThatRecovered(t *testing.T) { t.Errorf("Failed() = %d, want 0", got) } } + +// Every pass gets the refresher, retries included. Run feeds a failed item back +// through Download as the same struct it failed with, so without this a +// reference that expired mid-transfer would be replayed dead on every attempt. +func TestRunRefreshesOnEveryPassIncludingRetries(t *testing.T) { + var passes int + refresh := func(_ context.Context, it tgsource.Item) (tgsource.Item, error) { return it, nil } + + runFake(t, func(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o DownloadOptions) ([]Outcome, Stats, error) { + passes++ + if o.Refresh == nil { + t.Errorf("pass %d was given no refresher", passes) + } + return stageItems(map[int]bool{2: true})(ctx, seq, o) + }) + + o, staging, _ := runOpts(t, 0) + o.Refresh = refresh + + res, err := Run(t.Context(), seqOf(testItems(2, 10)), o) + if err != nil { + t.Fatalf("Run: %v", err) + } + if passes != downloadAttempts { + t.Fatalf("passes = %d, want %d — the retry passes are where a stale reference would be replayed", + passes, downloadAttempts) + } + if got := len(res.Failed()); got != 1 { + t.Errorf("Failed() = %d, want 1", got) + } + assertEmpty(t, staging) +} diff --git a/internal/tgsource/refresh.go b/internal/tgsource/refresh.go new file mode 100644 index 0000000..f8299fe --- /dev/null +++ b/internal/tgsource/refresh.go @@ -0,0 +1,129 @@ +package tgsource + +import ( + "context" + "errors" + "fmt" + + "github.com/gotd/td/telegram/peers" + "github.com/gotd/td/telegram/query" + "github.com/gotd/td/tg" + + "github.com/iyear/tdl/core/tmedia" + + "github.com/tiennm99dev/telegram-exporter/internal/naming" +) + +// ErrGone marks a message this run can no longer fetch: it was deleted, it no +// longer carries media, or the file behind it has been replaced by a different +// one. +// +// The distinction from an ordinary re-read failure is what the caller acts on. +// There is nothing here to retry and nothing a later pass would do differently, +// so one such message is left unarchived and the run carries on — whereas a +// re-read that failed because the connection did must count as a failure like +// any other, or a dead source stops looking like one. +var ErrGone = errors.New("message is gone") + +// Refresh re-reads one message and returns the item with a live file reference. +// +// A function type rather than an interface because there is one implementation +// and the download half only ever needs to call it; a func also lets a test +// supply a canned reference without a Telegram connection. +type Refresh func(context.Context, Item) (Item, error) + +// fetchMessage reads one message of a chat by id, reporting ErrGone if it is not +// there any more. It is the seam Refresher is built on, so the checks that +// decide whether a re-read describes the same file can be tested without a +// Telegram client. +type fetchMessage func(ctx context.Context, messageID int) (*tg.Message, error) + +// Refresher re-mints the file reference for an item from the chat it came from. +// +// The reference in Media.InputFileLoc is a short-lived token, and Telegram +// deliberately does not document how long it lives — the documented contract is +// the other way round: a client caches the reference *together with the source +// it came from*, and re-reads that source when the reference expires +// (core.telegram.org/api/file-references). This is that re-read. The source is +// the message, so reading the message again yields a document carrying a new +// reference. +func Refresher(api *tg.Client, peer peers.Peer) Refresh { + return refresher(historyFetch(api, peer.InputPeer())) +} + +func refresher(get fetchMessage) Refresh { + return func(ctx context.Context, it Item) (Item, error) { + msg, err := get(ctx, it.MessageID) + if err != nil { + if errors.Is(err, ErrGone) { + return Item{}, err + } + return Item{}, fmt.Errorf("re-read message %d: %w", it.MessageID, err) + } + + media, ok := tmedia.GetMedia(msg) + if !ok { + // The walk only yields messages that had media, so this means the + // message was edited into something else. + return Item{}, fmt.Errorf("%w: message %d no longer carries a file", + ErrGone, it.MessageID) + } + + // The re-read has to describe the same file, and both halves of that are + // checked because both are load-bearing. Name is what the remote index + // was diffed against and what the file will be written as, so a run that + // silently accepted a new one would archive different bytes under the + // old message's name. Size is what reserved the staging budget, what the + // finished download is measured against, and what the closing verify + // compares — a changed size would either fail that check every pass or + // be re-fetched forever as a size mismatch. + // + // A mismatch is therefore not a refresh failure but a different file, + // which this run has no name for and must leave alone. + if name := naming.For(it.DialogID, it.MessageID, media); name != it.Name || media.Size != it.Media.Size { + return Item{}, fmt.Errorf("%w: message %d now holds a different file (%q, %d bytes; was %q, %d bytes)", + ErrGone, it.MessageID, name, media.Size, it.Name, it.Media.Size) + } + + it.Media = media + return it, nil + } +} + +// historyFetch reads a single message out of a chat's history. +// +// This is a one-message GetHistory rather than channels.getMessages so that +// channels, groups and users need no branching, and it is written here rather +// than reusing tutil.GetSingleMessage because the classification is the point. +// tutil reports a deleted message three different ways — its own +// ErrMessageDeleted when an older message sits in the slot, "invalid message" +// when that message is a service one, and a wrapped nil when the page comes back +// empty — and only the first is distinguishable. This needs every one of them to +// be exactly ErrGone, because the caller retries anything else. +// +// OffsetID is exclusive, so id+1 asks for the message itself first. Telegram +// answers with the newest message at or below that offset, and gotd drops empty +// slots while paging (query/messages/iter.go), so a returned id lower than the +// one asked for means the message is not in the history any more. +func historyFetch(api *tg.Client, peer tg.InputPeerClass) fetchMessage { + return func(ctx context.Context, messageID int) (*tg.Message, error) { + it := query.NewQuery(api).Messages().GetHistory(peer). + OffsetID(messageID + 1).BatchSize(1).Iter() + + if !it.Next(ctx) { + // gotd ends the iterator with a nil error on an empty page, so a + // nil here is not success: it is "nothing at or below this id". + if err := it.Err(); err != nil { + return nil, err + } + return nil, fmt.Errorf("%w: no message at or before %d", ErrGone, messageID) + } + + msg, ok := it.Value().Msg.(*tg.Message) + if !ok || msg.ID != messageID { + return nil, fmt.Errorf("%w: message %d is no longer in the chat's history", + ErrGone, messageID) + } + return msg, nil + } +} diff --git a/internal/tgsource/refresh_test.go b/internal/tgsource/refresh_test.go new file mode 100644 index 0000000..e1ec2d2 --- /dev/null +++ b/internal/tgsource/refresh_test.go @@ -0,0 +1,167 @@ +package tgsource + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + + "github.com/gotd/td/tg" + + "github.com/iyear/tdl/core/tmedia" + + "github.com/tiennm99dev/telegram-exporter/internal/naming" +) + +// These drive refresher through its fetch seam rather than through a Telegram +// client, because what needs testing is not the round trip — it is the judgement +// afterwards. A re-read that describes a different file must be refused, and a +// re-read that failed because the message is gone must be distinguishable from +// one that failed because the connection did: the download half retries the +// second and skips the first. + +const testDialogID = 1234567890 + +// docMessage builds a message carrying one document, the shape tmedia extracts. +func docMessage(t *testing.T, id int, file string, size int64, ref string) *tg.Message { + t.Helper() + msg := &tg.Message{ID: id} + msg.SetMedia(&tg.MessageMediaDocument{ + Document: &tg.Document{ + ID: 1000000000000000001, + Size: size, + DCID: 2, + MimeType: "video/mp4", + FileReference: []byte(ref), + Attributes: []tg.DocumentAttributeClass{ + &tg.DocumentAttributeFilename{FileName: file}, + }, + }, + }) + return msg +} + +// walked is the item as the walk first saw it, carrying a reference that has +// since expired. +func walked(t *testing.T, id int, file string, size int64) Item { + t.Helper() + media, ok := tmedia.GetMedia(docMessage(t, id, file, size, "stale")) + if !ok { + t.Fatal("fixture message carries no media") + } + return Item{ + DialogID: testDialogID, + MessageID: id, + Name: naming.For(testDialogID, id, media), + Media: media, + } +} + +func TestRefresherSubstitutesALiveReference(t *testing.T) { + it := walked(t, 4242, "clip.mp4", 5000) + name := it.Name + + got, err := refresher(func(_ context.Context, id int) (*tg.Message, error) { + if id != 4242 { + t.Errorf("fetched message %d, want 4242", id) + } + return docMessage(t, 4242, "clip.mp4", 5000, "fresh"), nil + })(t.Context(), it) + if err != nil { + t.Fatalf("refresh: %v", err) + } + + loc, ok := got.Media.InputFileLoc.(*tg.InputDocumentFileLocation) + if !ok { + t.Fatalf("location is %T, want *tg.InputDocumentFileLocation", got.Media.InputFileLoc) + } + if string(loc.FileReference) != "fresh" { + t.Errorf("file reference = %q, want the re-read one", loc.FileReference) + } + // The name is the remote's key: it is what the run diffed against the index + // and what the file will be written as, so a refresh must never move it. + if got.Name != name { + t.Errorf("name = %q, want it unchanged at %q", got.Name, name) + } + if got.Media.Size != 5000 { + t.Errorf("size = %d, want 5000", got.Media.Size) + } +} + +// Every permanent reason a re-read cannot be used has to arrive as ErrGone, +// because that is the only signal that says "skip this one and keep going" +// rather than "the source is failing". +func TestRefresherReportsGoneForEveryPermanentCase(t *testing.T) { + it := walked(t, 4242, "clip.mp4", 5000) + + for _, tc := range []struct { + name string + get fetchMessage + want string + }{ + { + name: "message no longer in the history", + get: func(context.Context, int) (*tg.Message, error) { + return nil, fmt.Errorf("%w: no message at or before 4242", ErrGone) + }, + want: "4242", + }, + { + name: "message carries no media any more", + get: func(context.Context, int) (*tg.Message, error) { + return &tg.Message{ID: 4242, Message: "edited into text"}, nil + }, + want: "no longer carries a file", + }, + { + name: "a different file under the same message", + get: func(_ context.Context, _ int) (*tg.Message, error) { + return docMessage(t, 4242, "other.mp4", 5000, "fresh"), nil + }, + want: "different file", + }, + { + // Size is what reserved the staging budget, what the finished + // download is measured against and what the closing verify + // compares, so a same-name replacement of a different length is a + // different file too. + name: "same name, different length", + get: func(_ context.Context, _ int) (*tg.Message, error) { + return docMessage(t, 4242, "clip.mp4", 6000, "fresh"), nil + }, + want: "different file", + }, + } { + t.Run(tc.name, func(t *testing.T) { + _, err := refresher(tc.get)(t.Context(), it) + if !errors.Is(err, ErrGone) { + t.Fatalf("err = %v, want ErrGone", err) + } + if !strings.Contains(err.Error(), tc.want) { + t.Errorf("err = %v, want it to mention %q", err, tc.want) + } + }) + } +} + +// A re-read that failed for any other reason is not permanent and must not be +// dressed up as one: the caller counts these as download failures, retries them, +// and gives up on the source once enough of them land in a row. +func TestRefresherKeepsATransientFailureRetryable(t *testing.T) { + it := walked(t, 4242, "clip.mp4", 5000) + + _, err := refresher(func(context.Context, int) (*tg.Message, error) { + return nil, errors.New("dial telegram: connection refused") + })(t.Context(), it) + + if err == nil { + t.Fatal("a failed re-read reported success") + } + if errors.Is(err, ErrGone) { + t.Errorf("a connection failure was classified as ErrGone: %v", err) + } + if !strings.Contains(err.Error(), "4242") { + t.Errorf("err = %v, want it to name the message", err) + } +}