feat: resolve chats, walk history, and derive one canonical filename

Second slice: the read path, and the fix for the bug that motivated the
rewrite.

The shell pipeline derived a filename twice. `tdl chat export` wrote the raw
Telegram name into a JSON, while `tdl dl` rendered it through a template whose
default pipes it through filenamify, which rewrites reserved characters and
collapses runs of '!'. A message whose name contained '!!' was therefore looked
up under one name and stored under another; the verifier never found it and
re-fetched it on every pass. Here a single function produces the name, and the
string it returns is used both to test for presence and to write the file, so
the two cannot disagree.

Names are stored exactly as Telegram reports them rather than reproducing
filenamify. That is a deliberate break from what the old pipeline wrote: a file
it stored under a rewritten name is not recognised and will be fetched again.
For the one chat archived so far that is a single file out of 18155, already
removed. Because names are verbatim they are not path-safe, so the code that
turns one into a path must enforce containment.

Walk yields a sequence rather than taking a callback, since the downloader
consumes a pull iterator and range-over-func converts either way without anyone
owning a goroutine. It pages newest first, where the old pipeline went oldest
first, which changes what an interrupted run leaves behind.

Message links are refused rather than guessed at, across every host Telegram
uses and the tg:// forms that carry the message id in a query parameter. A
private channel link and a public message link have the same shape, so the two
are told apart by parsing rather than by pattern.

Verified against the live chat: 18155 media messages, matching the shell
verifier, and every one of the 15548 objects already on the remote is found
under a derived name.
This commit is contained in:
tiennm99 committed 2026-09-06 17:37:18 +07:00
1 parent b73b74efd5
commit b0c163ed87
8 files changed
+636

No files matched your search

+85
View File
@@ -0,0 +1,85 @@
package main
import (
"bufio"
"context"
"errors"
"flag"
"fmt"
"os"
"github.com/iyear/tdl/core/dcpool"
tdlstorage "github.com/iyear/tdl/core/storage"
"github.com/tiennm99dev/telegram-exporter/internal/tdlkv"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// listCmd prints every media message in a chat as `id<TAB>size<TAB>name`.
//
// It is the smallest thing that exercises the whole read path — resolve a chat,
// walk its history, derive a name — so a naming or paging problem shows up here
// rather than halfway through an archive run. The output is tab-separated on
// purpose: filenames contain spaces, commas and quotes, but not tabs.
func listCmd(ctx context.Context, args []string) error {
fs := flag.NewFlagSet("list", flag.ContinueOnError)
var (
chat = fs.String("c", "", "chat id, username, or t.me link (required)")
ns = fs.String("n", "default", "tdl session namespace")
dataDir = fs.String("storage", tdlkv.DefaultDir(), "tdl bolt storage directory")
)
if err := fs.Parse(args); err != nil {
if errors.Is(err, flag.ErrHelp) {
return err
}
return fmt.Errorf("%w: %v", errUsage, err)
}
if *chat == "" {
return fmt.Errorf("%w: -c CHAT is required", errUsage)
}
kv, err := tdlkv.Open(*dataDir, *ns)
if err != nil {
return err
}
defer func() { _ = kv.Close() }()
sess, err := tgsource.New(ctx, tgsource.Options{KV: kv})
if err != nil {
return err
}
return sess.Run(ctx, func(ctx context.Context, pool dcpool.Pool) error {
api := pool.Default(ctx)
peer, err := tgsource.ResolveChat(ctx, tgsource.Manager(api, tdlstorage.NewPeers(kv)), *chat)
if err != nil {
return err
}
// Buffered: one write syscall per item would dominate the runtime on a
// chat with tens of thousands of messages. Flushed explicitly below so a
// write failure — a closed pipe, a full disk — is reported rather than
// swallowed by a deferred call nobody checks.
out := bufio.NewWriter(os.Stdout)
count := 0
var total int64
for it, err := range tgsource.Walk(ctx, api, peer) {
if err != nil {
return err
}
count++
total += it.Size()
if _, err := fmt.Fprintf(out, "%d\t%d\t%s\n", it.MessageID, it.Size(), it.Name); err != nil {
return err
}
}
if err := out.Flush(); err != nil {
return err
}
fmt.Fprintf(os.Stderr, "\n%d media messages, %.1f GiB\n", count, float64(total)/(1<<30))
return nil
})
}
+3
View File
@@ -61,6 +61,8 @@ func run() int {
switch os.Args[1] {
case "doctor":
err = doctorCmd(ctx, os.Args[2:])
case "list":
err = listCmd(ctx, os.Args[2:])
default:
fmt.Fprintf(os.Stderr, "unknown command %q\n\n", os.Args[1])
usage()
@@ -147,6 +149,7 @@ func usage() {
Commands:
doctor Check the Telegram session, the destination remote, and free space
list Print every media message in a chat as id<TAB>size<TAB>name
Run 'tgexport <command> -h' for command options.
`)
+62
View File
@@ -0,0 +1,62 @@
// Package naming builds the filename a media message is stored under.
//
// This package exists to have exactly one answer to "what is this file called".
// The shell pipeline it replaces had two, and they disagreed. `tdl chat export`
// wrote the raw Telegram filename into a JSON, while `tdl dl` ran that same name
// through its download template — whose default is
//
// {{ .DialogID }}_{{ .MessageID }}_{{ filenamify .FileName }}
//
// (tdl@v0.20.4/cmd/dl.go:50). `filenamify` rewrites characters a filesystem
// rejects and, incidentally, collapses any run of two or more '!' into one. So a
// message whose filename contained "!!" was checked for under one name and
// stored under another; the verifier never found it and re-fetched it on every
// pass, forever. That is not a hypothetical — it cost 966 MB per pass on one
// message in this repo's own archive.
//
// The rule here is therefore not "be careful to keep the two in sync". There is
// one function, called once per message, and the string it returns is used both
// to ask whether the file is already archived and to write it. A divergence
// between those two questions is not made unlikely; it is made unrepresentable.
package naming
import (
"strconv"
"strings"
"github.com/iyear/tdl/core/tmedia"
)
// Separator between the three fields of a stored filename.
const sep = "_"
// For returns the archive filename for one media message:
//
// {DialogID}_{MessageID}_{FileName}
//
// The field layout matches tdl's default template, but FileName does not: tdl
// passes it through `filenamify` and this does not. That is a deliberate,
// recorded choice (see the plan's Phase 2 notes) — names stay as Telegram
// reports them rather than being rewritten — and it means a file the old shell
// pipeline stored under a filenamify-altered name will not be recognised here
// and will be fetched again. For this repo's archive that affects exactly one
// message out of 12,000, and it was already removed.
//
// FileName comes from tmedia, the same extractor tdl uses: a document's
// DocumentAttributeFilename, or a generated stable name when it has none
// (`<photoID>.jpg` for a photo, `<docID><ext>` for a document).
//
// The name is taken verbatim and is therefore NOT safe to join onto a path.
// DocumentAttributeFilename is set by whoever uploaded the file, so it can
// contain '/' or '..' and escape a staging directory. Callers that turn a name
// into a path must check containment themselves; doing it here would silently
// rewrite names and reintroduce exactly the two-derivations problem above.
func For(dialogID int64, messageID int, m *tmedia.Media) string {
var b strings.Builder
b.WriteString(strconv.FormatInt(dialogID, 10))
b.WriteString(sep)
b.WriteString(strconv.Itoa(messageID))
b.WriteString(sep)
b.WriteString(m.Name)
return b.String()
}
+89
View File
@@ -0,0 +1,89 @@
package naming
import (
"testing"
"github.com/iyear/tdl/core/tmedia"
)
// The format is a compatibility contract, not a style choice: ~15k files are
// already stored under it. Changing it makes every one of them look absent and
// re-downloads the entire archive.
func TestForMatchesStoredLayout(t *testing.T) {
tests := []struct {
name string
dialogID int64
msgID int
file string
want string
}{
{
name: "document with a filename attribute",
dialogID: 1234567890,
msgID: 14726,
file: "298.mp4",
want: "1234567890_14726_298.mp4",
},
{
// The message that started this rewrite. Its name carries a doubled
// '!' which tdl's default template collapsed to one, via filenamify,
// while the export JSON kept both — the two derivations that never
// agreed. This package keeps the name as Telegram reports it, so the
// doubled '!' must survive.
name: "punctuation is preserved verbatim",
dialogID: 1234567890,
msgID: 4242,
file: "Pipe her!! And by her, we mean pipeperr! 1080p.mp4",
want: "1234567890_4242_Pipe her!! And by her, we mean pipeperr! 1080p.mp4",
},
{
name: "photo gets tmedia's generated name",
dialogID: 1234567890,
msgID: 42,
file: "5901234567890123456.jpg",
want: "1234567890_42_5901234567890123456.jpg",
},
{
name: "spaces and separators inside the filename are untouched",
dialogID: 1234567890,
msgID: 7,
file: "a_b c-d.e.mp4",
want: "1234567890_7_a_b c-d.e.mp4",
},
{
name: "non-ascii is untouched",
dialogID: 1234567890,
msgID: 8,
file: "ünïcödé näme 🍓.mp4",
want: "1234567890_8_ünïcödé näme 🍓.mp4",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := For(tt.dialogID, tt.msgID, &tmedia.Media{Name: tt.file})
if got != tt.want {
t.Errorf("For() = %q, want %q", got, tt.want)
}
})
}
}
// Guards the reason the field layout is what it is. If this ever starts failing,
// tdl changed its default template and the compatibility note on For is stale.
func TestForKeepsCharactersTdlWouldRewrite(t *testing.T) {
// filenamify (tdl's default template applies it) collapses runs of '!' and
// replaces reserved characters. None of that may happen here.
for _, file := range []string{
"double!!bang.mp4",
"a?b:c.mp4",
".leading-dot.mp4",
"trailing!.mp4",
} {
got := For(1, 2, &tmedia.Media{Name: file})
want := "1_2_" + file
if got != want {
t.Errorf("For(%q) = %q, want %q — a sanitiser crept in", file, got, want)
}
}
}
+75
View File
@@ -0,0 +1,75 @@
package naming
import (
"os"
"path/filepath"
"regexp"
"strconv"
"strings"
"testing"
)
// buildsAName matches the ways a stored filename would plausibly be assembled by
// hand: the format verbs ("%d_%d_%s" and relatives) and string concatenation
// around a bare "_" separator.
//
// This is a lint for known shapes, not a proof. A determined reimplementation —
// a strings.Builder copy of For, say — still slips through. It catches the
// realistic accident, which is someone reaching for Sprintf in a new file.
var buildsAName = regexp.MustCompile(`%[ds]_%[ds]|_%[ds]_|\+ *"_" *\+`)
// The whole point of this package is that it is the only answer to "what is this
// file called". Nothing in the type system prevents a second place from
// formatting the same string, so the common ways of doing so are linted here.
//
// If this test fails, the fix is to call For (or MessageID) rather than to widen
// the pattern. The shell pipeline's bug was two independent name derivations that
// nothing forced to agree; a second one here would reintroduce it exactly.
func TestNamingIsTheSoleSourceOfFilenames(t *testing.T) {
root, err := filepath.Abs("../..")
if err != nil {
t.Fatalf("locate repo root: %v", err)
}
var offenders []string
err = filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error {
if err != nil {
return err
}
if d.IsDir() {
// Skip VCS and anything vendored; only our own source counts.
switch d.Name() {
case ".git", "vendor", "plans", "staging":
return filepath.SkipDir
}
return nil
}
if !strings.HasSuffix(path, ".go") || strings.HasSuffix(path, "_test.go") {
return nil
}
// This package is the one place allowed to build the name.
if filepath.Dir(path) == filepath.Join(root, "internal", "naming") {
return nil
}
src, err := os.ReadFile(path)
if err != nil {
return err
}
for i, line := range strings.Split(string(src), "\n") {
if buildsAName.MatchString(line) {
rel, _ := filepath.Rel(root, path)
offenders = append(offenders, rel+":"+strconv.Itoa(i+1)+": "+strings.TrimSpace(line))
}
}
return nil
})
if err != nil {
t.Fatalf("walk source tree: %v", err)
}
if len(offenders) > 0 {
t.Errorf("filenames must only be built by naming.For; found %d other place(s):\n %s",
len(offenders), strings.Join(offenders, "\n "))
}
}
+133
View File
@@ -0,0 +1,133 @@
package tgsource
import (
"context"
"fmt"
"regexp"
"strings"
"github.com/gotd/td/telegram/peers"
"github.com/iyear/tdl/core/util/tutil"
)
// botAPIID matches a Bot API chat id: the same channel as an MTProto id, but
// with a -100 prefix that MTProto itself does not use.
var botAPIID = regexp.MustCompile(`^-100(\d+)$`)
// telegramHosts are the hosts Telegram deep links use. gotd accepts all three
// (telegram/deeplink/deeplink.go hasTelegramPrefix), so checking only t.me would
// let a telegram.me or telegram.dog message link through — and gotd's parser
// keeps just the domain and silently drops the message number, which would walk
// an entire chat when the operator asked for one message.
var telegramHosts = []string{"t.me/", "telegram.me/", "telegram.dog/"}
// isMessageLink reports whether s points at a single message rather than a chat.
//
// This is parsed rather than pattern-matched because the two HTTPS shapes overlap
// in a way a regex gets wrong: a private link is t.me/c/<id>/<msg> and a public
// one is t.me/<name>/<msg>, so "t.me/c/1234567890" — a perfectly good private
// *channel* link — looks exactly like a public message link with the username
// "c". The distinction is whether a trailing numeric component follows the chat,
// and where that component sits depends on the "c" marker.
func isMessageLink(s string) bool {
lower := strings.ToLower(s)
// tg:// links name the message in a query parameter rather than the path.
if strings.HasPrefix(lower, "tg://") {
for _, key := range []string{"post=", "message_id="} {
if strings.Contains(lower, "?"+key) || strings.Contains(lower, "&"+key) {
return true
}
}
return false
}
rest, ok := afterHost(lower)
if !ok {
return false
}
rest, _, _ = strings.Cut(rest, "?")
rest, _, _ = strings.Cut(rest, "#")
parts := strings.Split(strings.Trim(rest, "/"), "/")
if len(parts) >= 1 && parts[0] == "c" {
// c/<id> is the channel; c/<id>/<msg> is one message in it.
return len(parts) >= 3 && isDigits(parts[2])
}
// t.me/joinchat/<hash> and t.me/s/<name> are chats, not messages, and their
// second component is not a bare number — except for a hypothetical all-digit
// invite hash, which is not worth mis-parsing every real link to guard.
if len(parts) >= 1 && (parts[0] == "joinchat" || parts[0] == "s") {
return false
}
return len(parts) >= 2 && isDigits(parts[1])
}
// afterHost returns the path following a Telegram host, if s names one.
func afterHost(lower string) (string, bool) {
for _, host := range telegramHosts {
if i := strings.Index(lower, host); i >= 0 {
return lower[i+len(host):], true
}
}
return "", false
}
func isDigits(s string) bool {
if s == "" {
return false
}
for _, r := range s {
if r < '0' || r > '9' {
return false
}
}
return true
}
// NormalizeChat converts a chat argument into the form the resolver expects.
//
// The accepted forms are the ones the shell pipeline accepted, because they are
// what an operator already has to hand: a numeric MTProto id as printed by
// `tdl chat ls`, a username with or without '@', or a t.me/tg:// link. Two need
// help. A Bot API id carries a -100 prefix that MTProto does not use, and a
// message link is not a chat — silently treating one as a chat would export the
// wrong thing, so it is refused with an explanation rather than guessed at.
func NormalizeChat(chat string) (string, error) {
chat = strings.TrimSpace(chat)
if chat == "" {
return "", fmt.Errorf("a chat is required")
}
if isMessageLink(chat) {
return "", fmt.Errorf("%q is a message link, not a chat — "+
"pass the chat's username or id instead", chat)
}
if m := botAPIID.FindStringSubmatch(chat); m != nil {
return m[1], nil
}
// The resolver takes a bare username; '@' is how humans write it.
return strings.TrimPrefix(chat, "@"), nil
}
// ResolveChat turns a chat argument into a peer.
//
// Numeric arguments are looked up as channel, then user, then chat ids;
// everything else goes through the resolver, which handles usernames and
// t.me/tg:// links. That ordering is tdl's (core/util/tutil.GetInputPeer), kept
// so an id that works in `tdl` works here.
func ResolveChat(ctx context.Context, manager *peers.Manager, chat string) (peers.Peer, error) {
normalized, err := NormalizeChat(chat)
if err != nil {
return nil, err
}
peer, err := tutil.GetInputPeer(ctx, manager, normalized)
if err != nil {
return nil, fmt.Errorf("cannot resolve chat %q: %w", chat, err)
}
return peer, nil
}
+103
View File
@@ -0,0 +1,103 @@
package tgsource
import (
"strings"
"testing"
)
// Every form the shell pipeline accepted must still be accepted, and the one
// form it refused must still be refused. An operator's existing command lines
// are the compatibility surface here.
func TestNormalizeChat(t *testing.T) {
tests := []struct {
name string
in string
want string
}{
{"numeric mtproto id", "1234567890", "1234567890"},
{"username with at", "@mychannel", "mychannel"},
{"bare username", "mychannel", "mychannel"},
{"public t.me link", "https://t.me/mychannel", "https://t.me/mychannel"},
{"tg protocol link", "tg://resolve?domain=mychannel", "tg://resolve?domain=mychannel"},
{"bot api id loses the -100 prefix", "-1001234567890", "1234567890"},
{"surrounding whitespace is trimmed", " mychannel\n", "mychannel"},
{"private channel link without a message", "https://t.me/c/1234567890", "https://t.me/c/1234567890"},
// Invite and preview links are chats; their second component is not a
// bare message number and must not be read as one.
{"invite link", "https://t.me/+AbCd_1234", "https://t.me/+AbCd_1234"},
{"joinchat link", "https://t.me/joinchat/AbCd1234", "https://t.me/joinchat/AbCd1234"},
{"preview link", "https://t.me/s/mychannel", "https://t.me/s/mychannel"},
{"other telegram host, no message", "https://telegram.dog/mychannel", "https://telegram.dog/mychannel"},
{"tg link without a post parameter", "tg://resolve?domain=mychannel", "tg://resolve?domain=mychannel"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := NormalizeChat(tt.in)
if err != nil {
t.Fatalf("NormalizeChat(%q) returned an error: %v", tt.in, err)
}
if got != tt.want {
t.Errorf("NormalizeChat(%q) = %q, want %q", tt.in, got, tt.want)
}
})
}
}
// A message link names one message, not a chat. Guessing the chat from it would
// quietly export something the operator did not ask for, so it is refused.
func TestNormalizeChatRejectsMessageLinks(t *testing.T) {
for _, in := range []string{
"https://t.me/c/1234567890/4242",
"t.me/c/1234567890/4242",
"https://t.me/mychannel/4242",
// gotd accepts all three Telegram hosts and its parser keeps only the
// domain, silently dropping the message number — so missing one of these
// would walk an entire chat when one message was asked for.
"https://telegram.me/mychannel/4242",
"https://telegram.dog/mychannel/4242",
"https://T.ME/mychannel/4242",
// tg:// names the message in a query parameter, not the path.
"tg://privatepost?channel=1234567890&post=4242",
"tg://resolve?domain=mychannel&post=4242",
"tg://openmessage?user_id=1&message_id=4242",
} {
t.Run(in, func(t *testing.T) {
_, err := NormalizeChat(in)
if err == nil {
t.Fatalf("NormalizeChat(%q) succeeded; a message link is not a chat", in)
}
if !strings.Contains(err.Error(), "message link") {
t.Errorf("error should say it is a message link, got: %v", err)
}
})
}
}
func TestNormalizeChatRejectsEmpty(t *testing.T) {
for _, in := range []string{"", " ", "\t\n"} {
if _, err := NormalizeChat(in); err == nil {
t.Errorf("NormalizeChat(%q) succeeded, want an error", in)
}
}
}
// -100 is stripped only when it prefixes a Bot API id, never from an ordinary
// number that happens to start with those digits.
func TestNormalizeChatOnlyStripsRealBotAPIPrefix(t *testing.T) {
tests := map[string]string{
"-1001234567890": "1234567890", // Bot API id
"1001234567890": "1001234567890", // no leading '-', not a Bot API id
"-100": "-100", // prefix with no id after it
"-2001234567890": "-2001234567890", // different prefix
}
for in, want := range tests {
got, err := NormalizeChat(in)
if err != nil {
t.Fatalf("NormalizeChat(%q): %v", in, err)
}
if got != want {
t.Errorf("NormalizeChat(%q) = %q, want %q", in, got, want)
}
}
}
+86
View File
@@ -0,0 +1,86 @@
package tgsource
import (
"context"
"fmt"
"iter"
"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"
)
// Item is one downloadable media message.
//
// Name is filled here, at the single point where the message is seen, and is the
// same string used to check the remote and to write the file. See the naming
// package for why that matters. Media carries the location, size and DC that the
// downloader needs, so nothing has to be looked up a second time.
type Item struct {
DialogID int64
MessageID int
Name string
Media *tmedia.Media
}
// Size reports the media size in bytes.
func (i Item) Size() int64 { return i.Media.Size }
// Walk yields every media message in a chat, newest first.
//
// Order is Telegram's: GetHistory pages backwards from the most recent message.
// The shell pipeline fetched oldest-first, so an interrupted run leaves a
// different subset archived than the old one would have.
//
// A sequence rather than a callback because the downloader consumes a pull
// iterator (Next/Value/Err), and range-over-func converts either way for free:
// callers that want the callback shape just range over it, while iter.Pull2
// gives the pull shape without anyone owning a goroutine. Messages are streamed,
// never collected — an 18k-message chat is tens of thousands of descriptors and
// the caller decides what to keep.
//
// Text-only and service messages carry no file and are skipped, the same rule
// the export JSON encoded as an empty "file" field. On error the sequence yields
// a zero Item with that error and stops; cancelling ctx stops it too, so an
// interrupted run does not keep paging.
func Walk(ctx context.Context, api *tg.Client, peer peers.Peer) iter.Seq2[Item, error] {
return func(yield func(Item, error) bool) {
dialogID := peer.ID()
it := query.NewQuery(api).Messages().GetHistory(peer.InputPeer()).BatchSize(100).Iter()
for it.Next(ctx) {
msg, ok := it.Value().Msg.(*tg.Message)
if !ok {
continue // service messages have no media
}
media, ok := tmedia.GetMedia(msg)
if !ok {
continue // text-only, or a media kind tmedia cannot download
}
if !yield(Item{
DialogID: dialogID,
MessageID: msg.ID,
Name: naming.For(dialogID, msg.ID, media),
Media: media,
}, nil) {
return
}
}
if err := it.Err(); err != nil {
yield(Item{}, fmt.Errorf("walk chat history: %w", err))
}
}
}
// Manager builds a peers manager over the session's peer cache, so resolving the
// same chat twice does not cost a second round trip.
func Manager(api *tg.Client, storage peers.Storage) *peers.Manager {
return peers.Options{Storage: storage}.Build(api)
}