fix: re-read each message for a live file reference before downloading it

The walk collects every item's file reference up front, and Telegram expires
those. The lifetime is undocumented; a run was seen to outlast an hour and fail
before two, which is less than a large chat spends downloading. So every fetch
past that point returned FILE_REFERENCE_EXPIRED, five in a row tripped the
breaker, and the run abandoned the rest of its todo list.

Telegram's contract is to cache a reference together with the source it came
from and re-read that source when it expires, so each message is re-read for a
live reference immediately before its own download. That is one round trip per
transfer and needs no guess at how long a reference lives. It runs after the
staging acquire rather than before: 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 reference all over again.

A re-read that describes a different file — different name or different length —
is refused, because archiving it under the name this run recorded would store
the wrong bytes. That case, a deleted message, and one that no longer holds a
file are all skipped: nothing can be archived for them, and one must not strand
the rest of the chat.

Any other re-read failure is recorded as a failed download instead, because that
is what it is. The connection that serves the re-read is the one that serves the
file, so these arrive exactly when transfers are failing too. An item that never
reaches a worker never reaches OnAdd or OnDone, so on the skip path it would
move no counter, feed no streak and produce no outcome to retry — a source
outage would walk the whole list one dead round trip at a time and still exit as
merely incomplete.

tutil.GetSingleMessage is not reused for the re-read: it reports a deleted
message three different ways, two of them wrapping a nil error, and only one is
distinguishable — so deletions would be misclassified as transient and spend the
breaker's streak.
This commit is contained in:
tiennm99 committed 2026-09-08 22:47:52 +07:00
1 parent efb29c9c39
commit b459b1f5e9
10 files changed
+773 -26

No files matched your search

+22
View File
@@ -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
+4
View File
@@ -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) {
+9
View File
@@ -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,
+258
View File
@@ -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)
}
})
}
}
+105 -26
View File
@@ -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 }
+7
View File
@@ -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,
+40
View File
@@ -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()
+32
View File
@@ -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)
}
+129
View File
@@ -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
}
}
+167
View File
@@ -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)
}
}