feat: index a remote and verify a chat against it in-process

Third slice: verify-export.sh, without the subprocess or the python.

One rclone listing builds an in-memory index, and the report is computed from
it. missing-ids.txt and gap.json are gone; so is every python3 heredoc.

Presence is answered from a whole name and never from a message id. The
id-keyed map exists only to tell "absent" apart from "absent, but a stale copy
under an older name is sitting there", and it stays unexported so nothing can
reach for it as an answer. That distinction is the bug this rewrite exists to
remove, so it is enforced by structure rather than by comment.

Names are checked for path containment before use. Storing them verbatim means
a filename chosen by whoever uploaded the file can contain a separator or a
parent reference, and tdl never had to care because its template rewrote those
away. Over-long names are refused for the same reason: the filesystem would
reject them at create time, and a file that can never be written would be
reported absent on every pass forever.

Objects are addressed by the path rclone knows them by, not by the basename
used for matching. The two differ once a remote has directory structure, and
deleting by basename would miss the object or remove a same-named one from the
root. Basenames appearing at more than one path make the snapshot ambiguous, so
they are reported rather than silently resolved.

Indexing runs at full depth with filters cleared. Inheriting RCLONE_MAX_DEPTH
or RCLONE_EXCLUDE would not fail, it would quietly report archived files as
absent and fetch them all again.

Deleting stale copies stays opt-in and confirmed; a non-interactive stdin
declines rather than proceeding. Filenames are quoted wherever they are
printed, so an embedded escape cannot redraw the list an operator approves.

Verified against the live remote: identical to verify-export.sh on the same
state — 18155 expected, 15548 present, 2607 absent, ids 9857-18013, exit 1.
This commit is contained in:
tiennm99 committed 2026-09-06 18:50:53 +07:00
1 parent b0c163ed87
commit 1393abdae2
10 files changed
+1242

No files matched your search

+3
View File
@@ -63,6 +63,8 @@ func run() int {
err = doctorCmd(ctx, os.Args[2:])
case "list":
err = listCmd(ctx, os.Args[2:])
case "verify":
err = verifyCmd(ctx, os.Args[2:])
default:
fmt.Fprintf(os.Stderr, "unknown command %q\n\n", os.Args[1])
usage()
@@ -150,6 +152,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
verify Report whether a chat is fully archived on a remote
Run 'tgexport <command> -h' for command options.
`)
+201
View File
@@ -0,0 +1,201 @@
package main
import (
"bufio"
"context"
"errors"
"flag"
"fmt"
"os"
"strings"
"github.com/iyear/tdl/core/dcpool"
tdlstorage "github.com/iyear/tdl/core/storage"
rclonefs "github.com/rclone/rclone/fs"
"github.com/rclone/rclone/fs/operations"
"github.com/tiennm99dev/telegram-exporter/internal/remote"
"github.com/tiennm99dev/telegram-exporter/internal/tdlkv"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
"github.com/tiennm99dev/telegram-exporter/internal/verify"
)
// verifyCmd reports whether a chat is fully archived on a remote.
//
// Exit 0 means complete, 1 means it ran and found work outstanding. Those are
// distinct on purpose: a driver needs to tell "nothing left to do" from "still
// incomplete" without parsing output.
func verifyCmd(ctx context.Context, args []string) error {
fs := flag.NewFlagSet("verify", flag.ContinueOnError)
var (
chat = fs.String("c", "", "chat id, username, or t.me link (required)")
remoteArg = fs.String("r", "", "rclone destination, e.g. pikpak:archive (required)")
ns = fs.String("n", "default", "tdl session namespace")
dataDir = fs.String("storage", tdlkv.DefaultDir(), "tdl bolt storage directory")
delStale = fs.Bool("delete-misnamed", false, "delete remote files stored under a superseded name")
assumeYes = fs.Bool("y", false, "do not prompt before deleting")
)
if err := fs.Parse(args); err != nil {
if errors.Is(err, flag.ErrHelp) {
return err
}
return fmt.Errorf("%w: %v", errUsage, err)
}
if *chat == "" || *remoteArg == "" {
return fmt.Errorf("%w: -c CHAT and -r REMOTE:PATH are both required", errUsage)
}
ctx, err := remote.Init(ctx, remote.DefaultTunables())
if err != nil {
return err
}
dst, err := remote.Resolve(ctx, *remoteArg)
if err != nil {
return err
}
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
}
var (
report verify.Report
idx *remote.Index
)
if err := 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
}
// The chat is walked first and the remote listed second, so the snapshot
// is never older than the wanted set. The reverse order could report a
// file absent that was uploaded while the walk was still running.
var items []tgsource.Item
for it, err := range tgsource.Walk(ctx, api, peer) {
if err != nil {
return err
}
items = append(items, it)
}
idx, err = remote.BuildIndex(ctx, dst, peer.ID())
if err != nil {
return err
}
fmt.Fprintf(os.Stderr, "indexed %d objects on %s\n", idx.Len(), dst.String())
if dup := idx.Collisions(); len(dup) > 0 {
// An ambiguous snapshot makes every verdict about these names
// unreliable, so it is reported rather than silently resolved.
fmt.Fprintf(os.Stderr, "warning: %d basename(s) appear at more than one path; "+
"verdicts for them may flip between runs:\n", len(dup))
for _, d := range dup[:min(5, len(dup))] {
fmt.Fprintf(os.Stderr, " %q\n", d)
}
}
fmt.Fprintln(os.Stderr)
report = verify.Check(items, idx)
return nil
}); err != nil {
return err
}
out := bufio.NewWriter(os.Stdout)
report.Write(out)
if err := out.Flush(); err != nil {
return err
}
if *delStale {
if err := deleteMisnamed(ctx, dst, idx, report, *assumeYes); err != nil {
return err
}
}
if !report.Complete() {
return fmt.Errorf("%w: %d file(s) still to fetch", errIncomplete, len(report.Todo()))
}
return nil
}
// deleteMisnamed removes stale copies left by an earlier naming scheme.
//
// Off by default and confirmed by default: this deletes data from the operator's
// remote, and a stale copy costs storage rather than correctness, so there is no
// hurry that justifies doing it unasked.
//
// Objects are addressed by the path the index recorded, not by the name matching
// used elsewhere. Those differ the moment a remote has directory structure, and
// deleting by basename would either miss the object or — worse, if the root
// happens to hold a same-named file — delete the wrong one.
func deleteMisnamed(ctx context.Context, dst rclonefs.Fs, idx *remote.Index, report verify.Report, assumeYes bool) error {
var targets []string
for _, m := range report.Misnamed {
targets = append(targets, m.Found...)
}
if len(targets) == 0 {
fmt.Fprintln(os.Stderr, "\nnothing to delete: no files stored under a superseded name")
return nil
}
// Filenames are chosen by whoever uploaded the file and may contain control
// characters or bidi marks, so they are quoted rather than printed raw: an
// embedded newline or escape sequence could otherwise redraw this list and
// have the operator approve something other than what they read.
fmt.Fprintf(os.Stderr, "\nabout to delete %d file(s) from %s:\n", len(targets), dst.String())
for _, t := range targets {
fmt.Fprintf(os.Stderr, " %q\n", t)
}
if !assumeYes {
// Anything that is not an explicit yes leaves the files alone, and that
// includes the read failing. Closed or non-interactive stdin — cron, a
// pipeline — therefore declines rather than proceeding, which is the
// safe direction for a delete.
fmt.Fprint(os.Stderr, "delete these? [y/N] ")
var answer string
_, _ = fmt.Scanln(&answer)
switch strings.ToLower(strings.TrimSpace(answer)) {
case "y", "yes":
default:
fmt.Fprintln(os.Stderr, "left alone")
return nil
}
}
// One failure must not strand the rest: the operator approved a set, so the
// whole set is attempted and the outcome reported as a count they can check
// against what they approved.
var errs []error
deleted := 0
for _, name := range targets {
path, ok := idx.PathOf(name)
if !ok {
errs = append(errs, fmt.Errorf("%q is no longer in the index", name))
continue
}
obj, err := dst.NewObject(ctx, path)
if err != nil {
errs = append(errs, fmt.Errorf("locate %q: %w", path, err))
continue
}
if err := operations.DeleteFile(ctx, obj); err != nil {
errs = append(errs, fmt.Errorf("delete %q: %w", path, err))
continue
}
deleted++
}
fmt.Fprintf(os.Stderr, "deleted %d of %d\n", deleted, len(targets))
return errors.Join(errs...)
}
+3
View File
@@ -87,3 +87,6 @@ func TestForKeepsCharactersTdlWouldRewrite(t *testing.T) {
}
}
}
// mediaNamed builds the only part of tmedia.Media these tests care about.
func mediaNamed(name string) *tmedia.Media { return &tmedia.Media{Name: name} }
+92
View File
@@ -0,0 +1,92 @@
package naming
import (
"fmt"
"path/filepath"
"strconv"
"strings"
)
// maxNameBytes is NAME_MAX on Linux: the longest single path component ext4 and
// friends accept. It is a byte count, not a rune count.
const maxNameBytes = 255
// Safe reports whether a stored name can be joined onto a directory path.
//
// Names are kept exactly as Telegram reports them, and the filename part comes
// from DocumentAttributeFilename — an unconstrained UTF-8 string chosen by
// whoever uploaded the file. So a name may contain '/' or be "..", and
// filepath.Join would happily resolve either outside the staging directory. tdl
// never had to think about this because its default template runs the name
// through filenamify, which rewrites separators and leading dots away; storing
// names verbatim moves that responsibility here.
//
// The rule is deliberately strict rather than corrective: a name that is not a
// single path element is rejected, not rewritten. Rewriting is what produced two
// disagreeing derivations in the first place, and a rejected file is a visible
// problem where a silently renamed one is not.
func Safe(name string) error {
switch {
case name == "":
return fmt.Errorf("empty filename")
case len(name) > maxNameBytes:
// The one rejection that is about a limit rather than an escape, and the
// one that matters most: os.Create returns ENAMETOOLONG past this, so an
// over-long name would download-fail forever while verify kept reporting
// it absent — a loop that never terminates. tdl could not hit this
// because filenamify truncates to 100 runes; storing names verbatim
// removes that cap, so the limit has to be checked instead.
return fmt.Errorf("filename is %d bytes, over the %d-byte limit: %q",
len(name), maxNameBytes, name)
case strings.ContainsRune(name, 0):
return fmt.Errorf("filename contains a NUL byte: %q", name)
case name == "." || name == "..":
return fmt.Errorf("filename is a directory reference: %q", name)
case filepath.IsAbs(name):
return fmt.Errorf("filename is an absolute path: %q", name)
case name != filepath.Base(name):
// Catches embedded separators, trailing slashes, and any ".." segment,
// since Base of all of those differs from the original.
return fmt.Errorf("filename is not a single path element: %q", name)
}
return nil
}
// SplitStored recovers the message id from a stored name, reporting false when
// the name does not have the expected shape or belongs to another dialog.
//
// Diagnostics only — telling "absent" apart from "present under a different
// name" when reporting on a remote. It must never decide that the wanted file is
// present: a file whose id matches but whose name does not is a different file.
// remote.Index keeps its id-keyed map unexported and its only id-to-name route
// excludes the wanted name, so an id cannot yield "the file you asked for is
// here" — though a caller that deliberately looks up a different name will of
// course get an answer about that name.
func SplitStored(dialogID int64, name string) (messageID int, ok bool) {
prefix := strconv.FormatInt(dialogID, 10) + sep
rest, found := strings.CutPrefix(name, prefix)
if !found {
return 0, false
}
idStr, _, found := strings.Cut(rest, sep)
if !found {
return 0, false
}
// Reject anything Atoi would accept but For would never emit: a sign, or
// leading zeros. The id has to be the exact text For wrote.
if idStr == "" || idStr[0] == '0' {
return 0, false
}
for _, r := range idStr {
if r < '0' || r > '9' {
return 0, false
}
}
// Only an overflow can fail here: the loop above rejected non-digits and a
// leading zero, so anything that parses is already >= 1.
id, err := strconv.Atoi(idStr)
if err != nil {
return 0, false
}
return id, true
}
+147
View File
@@ -0,0 +1,147 @@
package naming
import (
"path/filepath"
"strings"
"testing"
)
// Names come from DocumentAttributeFilename, which whoever uploaded the file
// chose. Anything that is not a single path element must be refused before it
// reaches filepath.Join.
func TestSafeRejectsNamesThatEscapeADirectory(t *testing.T) {
bad := []struct {
name string
want string
}{
{"", "empty"},
{".", "directory reference"},
{"..", "directory reference"},
{"../escape.mp4", "single path element"},
{"../../../.config/rclone/rclone.conf", "single path element"},
{"sub/dir.mp4", "single path element"},
{"trailing/", "single path element"},
{"/etc/passwd", "absolute path"},
{"/", "absolute path"},
{"nul\x00byte.mp4", "NUL byte"},
}
for _, tt := range bad {
t.Run(tt.name, func(t *testing.T) {
err := Safe(tt.name)
if err == nil {
t.Fatalf("Safe(%q) = nil; this name escapes or breaks a path", tt.name)
}
if !strings.Contains(err.Error(), tt.want) {
t.Errorf("Safe(%q) error = %v, want it to mention %q", tt.name, err, tt.want)
}
})
}
}
// Ordinary media names, including awkward but legal ones, must pass — a
// containment check that rejects real files is just a different outage.
func TestSafeAcceptsRealNames(t *testing.T) {
for _, name := range []string{
"1234567890_4242_Pipe her!! And by her, we mean pipeperr! 1080p.mp4",
"1234567890_14726_298.mp4",
"1234567890_8_ünïcödé näme 🍓.mp4",
"1234567890_42_a?b:c.mp4", // reserved on Windows, fine here
"1234567890_7_...dots.mp4", // leading dots inside the field, not the name
"1234567890_9_-dash.mp4",
"1234567890_10_ leading-space.mp4",
} {
if err := Safe(name); err != nil {
t.Errorf("Safe(%q) = %v, want nil", name, err)
}
}
}
// The property that actually matters: a name Safe accepts cannot, once joined,
// resolve outside the directory it was joined to.
func TestSafeNamesStayInsideTheStagingDirectory(t *testing.T) {
const staging = "/var/tmp/staging"
for _, name := range []string{
"1234567890_1_ok.mp4",
"1234567890_2_..dots.mp4",
"1234567890_3_a..b.mp4",
} {
if err := Safe(name); err != nil {
t.Fatalf("Safe(%q) = %v, want nil", name, err)
}
joined := filepath.Clean(filepath.Join(staging, name))
if filepath.Dir(joined) != staging {
t.Errorf("Join(%q, %q) = %q, which leaves the staging directory", staging, name, joined)
}
}
}
func TestSplitStoredRoundTripsWhatForProduces(t *testing.T) {
const dialog = int64(1234567890)
for _, msgID := range []int{1, 42, 4242, 4246} {
name := For(dialog, msgID, mediaNamed("a_b!!.mp4"))
got, ok := SplitStored(dialog, name)
if !ok {
t.Errorf("SplitStored(%q) reported no match", name)
continue
}
if got != msgID {
t.Errorf("SplitStored(%q) = %d, want %d", name, got, msgID)
}
}
}
func TestSplitStoredRejectsForeignNames(t *testing.T) {
const dialog = int64(1234567890)
for _, name := range []string{
"999_42_other-dialog.mp4", // different dialog
"1234567890_notanumber_.mp4", // id is not a number
"1234567890_0_zero.mp4", // ids start at 1
"1234567890_042_pad.mp4", // For never emits leading zeros
"1234567890_+42_sign.mp4", // nor a sign
"1234567890_-5_neg.mp4",
"1234567890_42", // no filename field
"1234567890", // no id field
"",
} {
if id, ok := SplitStored(dialog, name); ok {
t.Errorf("SplitStored(%q) = %d, true; want no match", name, id)
}
}
}
// A dialog id that is a prefix of another must not match it.
func TestSplitStoredDoesNotMatchPrefixOverlap(t *testing.T) {
if id, ok := SplitStored(123456789, "1234567890_42_file.mp4"); ok {
t.Errorf("SplitStored matched a longer dialog id, got %d", id)
}
}
// An over-long name is the one input that can hang a drive-until-complete loop:
// os.Create rejects it, so the download fails forever while verify keeps
// reporting it absent. It must be refused up front, not discovered per pass.
func TestSafeRejectsNamesOverTheFilesystemLimit(t *testing.T) {
prefix := "1234567890_42_"
fill := 255 - len(prefix)
atLimit := prefix + strings.Repeat("a", fill)
if err := Safe(atLimit); err != nil {
t.Errorf("Safe(%d bytes) = %v, want nil at exactly the limit", len(atLimit), err)
}
overLimit := prefix + strings.Repeat("a", fill+1)
err := Safe(overLimit)
if err == nil {
t.Fatalf("Safe(%d bytes) = nil, want an error past the limit", len(overLimit))
}
if !strings.Contains(err.Error(), "over the 255-byte limit") {
t.Errorf("error should name the limit, got: %v", err)
}
// The limit is bytes, not runes: multi-byte names hit it sooner.
multibyte := prefix + strings.Repeat("é", 130) // 260 bytes of payload
if err := Safe(multibyte); err == nil {
t.Errorf("Safe(%d bytes, %d runes) = nil; the limit must count bytes",
len(multibyte), len([]rune(multibyte)))
}
}
+131
View File
@@ -0,0 +1,131 @@
// Package pipeline drives the download half of an archive run.
package pipeline
import (
"context"
"fmt"
"io"
"iter"
"os"
"path/filepath"
"github.com/gotd/td/tg"
"github.com/iyear/tdl/core/downloader"
"github.com/tiennm99dev/telegram-exporter/internal/naming"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// partSuffix marks a download that is still in flight.
//
// Deliberately not tdl's ".tmp": nothing here excludes by extension any more,
// because upload is triggered by a download returning rather than by a filter
// over a directory. A distinct suffix just keeps a staging directory shared with
// a legacy tdl run unambiguous during the cutover.
const partSuffix = ".part"
// elem adapts one media item to the downloader's element interface.
type elem struct {
item tgsource.Item
file *os.File
takeout bool
}
func (e *elem) File() downloader.File { return mediaFile{e.item} }
func (e *elem) To() io.WriterAt { return e.file }
func (e *elem) AsTakeout() bool { return e.takeout }
// mediaFile exposes what the downloader needs to locate the bytes. All three
// values come straight from tmedia, so nothing is looked up a second time.
type mediaFile struct{ item tgsource.Item }
func (f mediaFile) Location() tg.InputFileLocationClass { return f.item.Media.InputFileLoc }
func (f mediaFile) Size() int64 { return f.item.Media.Size }
func (f mediaFile) DC() int { return f.item.Media.DC }
// partPath and finalPath are where an item is written and where it lands.
func partPath(staging string, it tgsource.Item) string {
return filepath.Join(staging, it.Name+partSuffix)
}
func finalPath(staging string, it tgsource.Item) string {
return filepath.Join(staging, it.Name)
}
// elemIter turns the item sequence into the pull iterator the downloader wants,
// opening each destination file as it goes.
//
// The downloader consumes Next/Value/Err; Walk produces an iter.Seq2. iter.Pull2
// bridges them without this code owning a goroutine or a channel, which is why
// Walk returns a sequence in the first place.
type elemIter struct {
next func() (tgsource.Item, error, bool)
stop func()
staging string
takeout bool
current *elem
err error
// opened records every file handle so a run can close them all. The
// downloader never closes what To() hands it, and a leak here is thousands
// of descriptors on a full archive run.
opened []*os.File
}
func newElemIter(seq iter.Seq2[tgsource.Item, error], staging string, takeout bool) *elemIter {
next, stop := iter.Pull2(seq)
return &elemIter{next: next, stop: stop, staging: staging, takeout: takeout}
}
func (i *elemIter) Next(ctx context.Context) bool {
if i.err != nil {
return false
}
if err := ctx.Err(); err != nil {
i.err = err
return false
}
item, err, ok := i.next()
if !ok {
return false
}
if err != nil {
i.err = err
return false
}
// A name that cannot be written is refused here rather than left to
// os.Create: the error names the message, and the run continues instead of
// failing on a path that could never have worked.
if err := naming.Safe(item.Name); err != nil {
i.err = fmt.Errorf("message %d: %w", item.MessageID, err)
return false
}
f, err := os.OpenFile(partPath(i.staging, item), os.O_CREATE|os.O_RDWR, 0o600)
if err != nil {
i.err = fmt.Errorf("open destination for message %d: %w", item.MessageID, err)
return false
}
i.opened = append(i.opened, f)
i.current = &elem{item: item, file: f, takeout: i.takeout}
return true
}
func (i *elemIter) Value() downloader.Elem { return i.current }
func (i *elemIter) Err() error { return i.err }
// Close releases the pull iterator and every file the walk opened.
func (i *elemIter) Close() error {
i.stop()
var firstErr error
for _, f := range i.opened {
if err := f.Close(); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
}
+134
View File
@@ -0,0 +1,134 @@
package remote
import (
"context"
"fmt"
"path"
"slices"
"github.com/rclone/rclone/fs"
"github.com/rclone/rclone/fs/filter"
"github.com/rclone/rclone/fs/operations"
"github.com/tiennm99dev/telegram-exporter/internal/naming"
)
// object is what the index remembers about one stored file.
//
// Path is kept alongside the basename because acting on an object — deleting a
// stale copy, say — needs the path rclone knows it by, while matching needs the
// basename. Conflating the two makes a delete address the wrong file, or no file
// at all, whenever a remote has any directory structure.
type object struct {
path string
size int64
}
// Index is a snapshot of what a remote holds, keyed by filename.
//
// It answers one question — "is this exact name present, and how big is it?" —
// and it answers it from the whole name, never from a message id. That is the
// point: a file whose id matches but whose name does not is a different file,
// and treating it as present is precisely the bug this rewrite exists to remove.
//
// The id-keyed map below exists only to tell "absent" apart from "absent, but
// something else is stored under this message's id" when reporting. It is
// unexported, and the only exported route from an id to a name is
// StoredUnderOtherNames, which by construction excludes the wanted name — so no
// caller can get "the file you asked for is present" out of an id.
type Index struct {
byName map[string]object
byID map[int][]string
collisions []string
}
// BuildIndex lists a remote once and indexes it.
//
// Object paths are reduced to their basename, so a remote written with a
// subdirectory layout matches the same way a flat one does — the shell verifier
// did this too, and existing archives rely on it.
//
// Listing runs with depth and filters neutralised. rclone's ListFn otherwise
// inherits whatever RCLONE_MAX_DEPTH or RCLONE_EXCLUDE happen to be set to, and
// a narrowed listing here does not fail — it silently reports archived files as
// absent and re-downloads every one of them. The transfer tunables in Init are
// deliberately env-overridable; this is not.
func BuildIndex(ctx context.Context, f fs.Fs, dialogID int64) (*Index, error) {
ctx, ci := fs.AddConfig(ctx)
ci.MaxDepth = -1
unfiltered, err := filter.NewFilter(nil)
if err != nil {
return nil, fmt.Errorf("build an empty filter: %w", err)
}
ctx = filter.ReplaceConfig(ctx, unfiltered)
idx := &Index{
byName: make(map[string]object),
byID: make(map[int][]string),
}
// ListFn is documented not to call fn concurrently, so the maps need no lock.
if err := operations.ListFn(ctx, f, func(o fs.Object) {
name := path.Base(o.Remote())
if _, seen := idx.byName[name]; seen {
// Two objects in different directories sharing a basename. Listing
// order is not guaranteed, so silently keeping one would make the
// verdict flip between runs — a zero-byte copy and a complete one
// would alternate. Keep the first and report the ambiguity instead.
if !slices.Contains(idx.collisions, name) {
idx.collisions = append(idx.collisions, name)
}
return
}
idx.byName[name] = object{path: o.Remote(), size: o.Size()}
if id, ok := naming.SplitStored(dialogID, name); ok {
idx.byID[id] = append(idx.byID[id], name)
}
}); err != nil {
return nil, fmt.Errorf("list %s: %w", f.String(), err)
}
return idx, nil
}
// Lookup reports the size stored under an exact name.
func (i *Index) Lookup(name string) (size int64, ok bool) {
o, ok := i.byName[name]
return o.size, ok
}
// PathOf returns the remote path an indexed name was found at, which is what
// rclone needs to act on the object. It differs from the name whenever the
// remote has directory structure.
func (i *Index) PathOf(name string) (string, bool) {
o, ok := i.byName[name]
return o.path, ok
}
// Len reports how many distinct names the remote held when the snapshot was
// taken. Objects dropped as basename collisions are not counted.
func (i *Index) Len() int { return len(i.byName) }
// Collisions lists basenames that appeared at more than one path. A non-empty
// result means the snapshot is ambiguous and any verdict about those names is
// unreliable, so callers should surface it rather than ignore it.
func (i *Index) Collisions() []string { return i.collisions }
// StoredUnderOtherNames lists names present for a message id that are not the
// wanted name.
//
// Diagnostics only. A non-empty result never means the file is archived — it
// means a stale copy from an earlier naming scheme is sitting there and will
// still be sitting there after the re-download, which is why the report has to
// surface it rather than quietly counting it.
func (i *Index) StoredUnderOtherNames(messageID int, wanted string) []string {
var others []string
for _, n := range i.byID[messageID] {
if n != wanted {
others = append(others, n)
}
}
return others
}
+174
View File
@@ -0,0 +1,174 @@
package remote
import (
"context"
"os"
"path/filepath"
"testing"
_ "github.com/rclone/rclone/backend/local"
"github.com/rclone/rclone/fs"
)
const testDialog = int64(1234567890)
// localIndex builds an index over a temp directory using rclone's local
// backend, so BuildIndex is exercised through the same listing path a real
// remote uses rather than through a stub.
func localIndex(t *testing.T, files map[string]int) *Index {
t.Helper()
dir := t.TempDir()
for name, size := range files {
full := filepath.Join(dir, name)
if err := os.MkdirAll(filepath.Dir(full), 0o755); err != nil {
t.Fatalf("mkdir for %q: %v", name, err)
}
if err := os.WriteFile(full, make([]byte, size), 0o600); err != nil {
t.Fatalf("write %q: %v", name, err)
}
}
ctx := context.Background()
f, err := fs.NewFs(ctx, dir)
if err != nil {
t.Fatalf("open local fs: %v", err)
}
idx, err := BuildIndex(ctx, f, testDialog)
if err != nil {
t.Fatalf("BuildIndex: %v", err)
}
return idx
}
func TestLookupMatchesWholeNamesOnly(t *testing.T) {
idx := localIndex(t, map[string]int{
"1234567890_4242_Pipe her! And by her, we mean pipeperr! 1080p.mp4": 4096,
"1234567890_14726_298.mp4": 2048,
})
if got, ok := idx.Lookup("1234567890_14726_298.mp4"); !ok || got != 2048 {
t.Errorf("Lookup(exact) = %d, %v; want 2048, true", got, ok)
}
// The doubled '!' is the name Telegram reports; the remote holds the
// collapsed one that tdl's filenamify template wrote. Those are different
// files as far as this index is concerned, and that is the whole policy.
wanted := "1234567890_4242_Pipe her!! And by her, we mean pipeperr! 1080p.mp4"
if _, ok := idx.Lookup(wanted); ok {
t.Error("Lookup matched a near-miss name; presence must require an exact match")
}
}
// A remote written with a subdirectory layout has to match a flat one, because
// existing archives were written both ways.
func TestBuildIndexReducesPathsToBasename(t *testing.T) {
idx := localIndex(t, map[string]int{
"nested/dir/1234567890_42_deep.mp4": 512,
})
if got, ok := idx.Lookup("1234567890_42_deep.mp4"); !ok || got != 512 {
t.Errorf("Lookup after basename reduction = %d, %v; want 512, true", got, ok)
}
}
func TestStoredUnderOtherNamesFindsStaleCopies(t *testing.T) {
stale := "1234567890_4242_Pipe her! And by her, we mean pipeperr! 1080p.mp4"
idx := localIndex(t, map[string]int{stale: 4096})
wanted := "1234567890_4242_Pipe her!! And by her, we mean pipeperr! 1080p.mp4"
others := idx.StoredUnderOtherNames(4242, wanted)
if len(others) != 1 || others[0] != stale {
t.Fatalf("StoredUnderOtherNames = %v, want [%q]", others, stale)
}
// The wanted name itself is never reported as an "other" name.
idx2 := localIndex(t, map[string]int{wanted: 4096})
if others := idx2.StoredUnderOtherNames(4242, wanted); len(others) != 0 {
t.Errorf("StoredUnderOtherNames = %v, want empty when the wanted name is present", others)
}
}
// Objects belonging to a different dialog, or not matching the stored layout at
// all, must not be indexed by id — otherwise an unrelated file could be reported
// as a stale copy of a message.
func TestStoredUnderOtherNamesIgnoresForeignObjects(t *testing.T) {
idx := localIndex(t, map[string]int{
"999999_4242_other-dialog.mp4": 100,
"not-a-tdl-name.mp4": 100,
})
if others := idx.StoredUnderOtherNames(4242, "1234567890_4242_x.mp4"); len(others) != 0 {
t.Errorf("StoredUnderOtherNames = %v, want empty", others)
}
}
func TestIndexLenCountsEveryObject(t *testing.T) {
idx := localIndex(t, map[string]int{
"1234567890_1_a.mp4": 1,
"1234567890_2_b.mp4": 1,
"unrelated.txt": 1,
})
if idx.Len() != 3 {
t.Errorf("Len() = %d, want 3", idx.Len())
}
}
func TestBuildIndexOnEmptyRemote(t *testing.T) {
idx := localIndex(t, nil)
if idx.Len() != 0 {
t.Errorf("Len() = %d, want 0", idx.Len())
}
if _, ok := idx.Lookup("anything"); ok {
t.Error("Lookup on an empty index reported a hit")
}
}
// Two objects in different directories can share a basename. Listing order is
// not guaranteed, so silently keeping one would make the verdict flip between
// runs; the ambiguity has to be reported instead.
func TestBuildIndexReportsBasenameCollisions(t *testing.T) {
idx := localIndex(t, map[string]int{
"a/1234567890_42_same.mp4": 100,
"b/1234567890_42_same.mp4": 0,
})
dup := idx.Collisions()
if len(dup) != 1 || dup[0] != "1234567890_42_same.mp4" {
t.Fatalf("Collisions() = %v, want the shared basename reported", dup)
}
if idx.Len() != 1 {
t.Errorf("Len() = %d, want 1 — a dropped collision must not be counted", idx.Len())
}
// The id map must not gain a duplicate entry either, or a delete would try
// the same name twice and fail the second time.
others := idx.StoredUnderOtherNames(42, "1234567890_42_wanted.mp4")
if len(others) != 1 {
t.Errorf("StoredUnderOtherNames = %v, want one entry, not a duplicate", others)
}
}
// Acting on an object needs the path rclone knows it by, which differs from the
// basename used for matching whenever the remote has directory structure.
// Deleting by basename would miss the object, or hit the wrong one.
func TestPathOfReturnsTheFullRemotePath(t *testing.T) {
idx := localIndex(t, map[string]int{"nested/dir/1234567890_42_deep.mp4": 512})
const name = "1234567890_42_deep.mp4"
if _, ok := idx.Lookup(name); !ok {
t.Fatalf("Lookup(%q) missed; matching is by basename", name)
}
got, ok := idx.PathOf(name)
if !ok {
t.Fatalf("PathOf(%q) reported no match", name)
}
if want := "nested/dir/" + name; got != want {
t.Errorf("PathOf(%q) = %q, want %q — deleting by basename would target the wrong path", name, got, want)
}
}
func TestPathOfMissesUnknownNames(t *testing.T) {
idx := localIndex(t, nil)
if p, ok := idx.PathOf("absent.mp4"); ok {
t.Errorf("PathOf(absent) = %q, true; want no match", p)
}
}
+176
View File
@@ -0,0 +1,176 @@
// Package verify answers whether a chat is fully archived on a remote.
//
// It replaces verify-export.sh and keeps that script's size judgements, which
// were arrived at by watching real failures rather than by taste.
//
// One judgement is deliberately not carried over. The script also counted files
// sitting in the staging directory as present, because download and upload were
// separate processes and a file could be finished locally but not yet uploaded
// for a whole sync interval. Here a single process owns both legs, so that state
// is not one a verify can meaningfully observe — except after an interrupted
// run, where staging may hold completed files. Phase 5 owns staging and decides
// whether to credit it; until then a verify after an interrupt may report files
// absent that are on local disk, and re-fetch them.
package verify
import (
"fmt"
"io"
"sort"
"github.com/tiennm99dev/telegram-exporter/internal/naming"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
// Index is the part of a remote snapshot verification needs.
//
// An interface rather than *remote.Index so the two questions stay separable:
// presence is answered from a whole name, and the id-keyed lookup is explicitly
// a different method with a name that says it is not an answer. It also lets the
// report be tested without a remote.
type Index interface {
// Lookup reports the size stored under an exact name.
Lookup(name string) (size int64, ok bool)
// StoredUnderOtherNames lists names present for a message id that are not
// the wanted name. Diagnostics only; never a presence answer.
StoredUnderOtherNames(messageID int, wanted string) []string
}
// tinyThreshold is the size below which a present file is reported for a human
// to look at but still trusted. Some real media genuinely is this small, so
// treating it as damaged would re-download it forever.
const tinyThreshold = 1024
// Misnamed is a wanted file that is absent, while some other file is stored
// under the same message id.
type Misnamed struct {
MessageID int
Wanted string
Found []string
}
// Unsafe is a wanted file whose name cannot be written to a path.
type Unsafe struct {
MessageID int
Name string
Reason error
}
// Tiny is a present file small enough to be worth a look.
type Tiny struct {
MessageID int
Name string
Size int64
}
// Report is the outcome of comparing a chat against a remote.
type Report struct {
Expected int // media messages in the chat
Present int // present, non-empty
Bytes int64 // total size of everything expected
Absent []int // not on the remote under the wanted name
ZeroByte []int // present but empty
Misnamed []Misnamed
Unsafe []Unsafe
Tiny []Tiny
}
// Todo lists the message ids needing another fetch, in ascending order.
//
// Zero-byte files are included: rclone overwrites a size-mismatched destination,
// so simply fetching again repairs them.
func (r Report) Todo() []int {
todo := make([]int, 0, len(r.Absent)+len(r.ZeroByte))
todo = append(todo, r.Absent...)
todo = append(todo, r.ZeroByte...)
sort.Ints(todo)
return todo
}
// Complete reports whether every expected file is present and non-empty.
func (r Report) Complete() bool { return len(r.Todo()) == 0 }
// Check compares the wanted items against an index of the remote.
//
// Matching is on the whole name. A file stored under any other name is not the
// file that was asked for, however close it looks — that is a deliberate policy,
// and the near-misses are collected into Misnamed rather than being quietly
// accepted, because the re-download lands beside them and both copies stay.
func Check(items []tgsource.Item, idx Index) Report {
r := Report{Expected: len(items)}
for _, it := range items {
r.Bytes += it.Size()
if err := naming.Safe(it.Name); err != nil {
// Never counted present: this name cannot be written anywhere safe,
// so no correctly-behaving run could have archived it.
r.Unsafe = append(r.Unsafe, Unsafe{MessageID: it.MessageID, Name: it.Name, Reason: err})
r.Absent = append(r.Absent, it.MessageID)
continue
}
size, ok := idx.Lookup(it.Name)
switch {
case !ok:
r.Absent = append(r.Absent, it.MessageID)
if others := idx.StoredUnderOtherNames(it.MessageID, it.Name); len(others) > 0 {
r.Misnamed = append(r.Misnamed, Misnamed{
MessageID: it.MessageID, Wanted: it.Name, Found: others,
})
}
case size == 0:
r.ZeroByte = append(r.ZeroByte, it.MessageID)
default:
r.Present++
if size < tinyThreshold {
r.Tiny = append(r.Tiny, Tiny{MessageID: it.MessageID, Name: it.Name, Size: size})
}
}
}
return r
}
// Write renders a report in the shape verify-export.sh printed, so the numbers
// stay comparable across the cutover.
func (r Report) Write(w io.Writer) {
fmt.Fprintf(w, "media expected : %d (%.1f GiB)\n", r.Expected, float64(r.Bytes)/(1<<30))
fmt.Fprintf(w, "present and intact : %d\n", r.Present)
fmt.Fprintf(w, " absent : %d\n", len(r.Absent))
fmt.Fprintf(w, " zero-byte : %d\n", len(r.ZeroByte))
if len(r.Tiny) > 0 {
fmt.Fprintf(w, " under 1KiB (check, not retried): %d\n", len(r.Tiny))
for _, t := range r.Tiny[:min(5, len(r.Tiny))] {
fmt.Fprintf(w, " id %d %d B %q\n", t.MessageID, t.Size, t.Name)
}
}
if len(r.Unsafe) > 0 {
fmt.Fprintf(w, "\nunsafe filenames : %d\n", len(r.Unsafe))
fmt.Fprintf(w, " these cannot be written to a path and are never fetched:\n")
for _, u := range r.Unsafe {
fmt.Fprintf(w, " id %d %v\n", u.MessageID, u.Reason)
}
}
if len(r.Misnamed) > 0 {
fmt.Fprintf(w, "\nstored under a different name : %d\n", len(r.Misnamed))
fmt.Fprintf(w, " counted as absent and fetched again; delete the stale copies so the\n")
fmt.Fprintf(w, " re-download does not leave two files for the same message:\n")
for _, m := range r.Misnamed {
fmt.Fprintf(w, " id %d\n wanted: %q\n", m.MessageID, m.Wanted)
for _, f := range m.Found {
fmt.Fprintf(w, " remote: %q\n", f)
}
}
}
if todo := r.Todo(); len(todo) > 0 {
fmt.Fprintf(w, "\nneeds another pass : %d (ids %d–%d)\n", len(todo), todo[0], todo[len(todo)-1])
return
}
fmt.Fprintf(w, "\nCOMPLETE: every media message is present and non-empty.\n")
}
+181
View File
@@ -0,0 +1,181 @@
package verify
import (
"slices"
"strings"
"testing"
"github.com/iyear/tdl/core/tmedia"
"github.com/tiennm99dev/telegram-exporter/internal/naming"
"github.com/tiennm99dev/telegram-exporter/internal/tgsource"
)
const dialog = int64(1234567890)
// fakeIndex is a remote snapshot expressed directly, so report logic is tested
// without a network or a filesystem.
type fakeIndex map[string]int64
func (f fakeIndex) Lookup(name string) (int64, bool) {
size, ok := f[name]
return size, ok
}
func (f fakeIndex) StoredUnderOtherNames(messageID int, wanted string) []string {
var others []string
for name := range f {
if name == wanted {
continue
}
if id, ok := naming.SplitStored(dialog, name); ok && id == messageID {
others = append(others, name)
}
}
slices.Sort(others)
return others
}
func item(msgID int, file string, size int64) tgsource.Item {
m := &tmedia.Media{Name: file, Size: size}
return tgsource.Item{
DialogID: dialog,
MessageID: msgID,
Name: naming.For(dialog, msgID, m),
Media: m,
}
}
func TestCheckClassifiesEveryOutcome(t *testing.T) {
items := []tgsource.Item{
item(1, "present.mp4", 5000),
item(2, "empty.mp4", 5000),
item(3, "absent.mp4", 5000),
item(4, "tiny.jpg", 500),
}
idx := fakeIndex{
"1234567890_1_present.mp4": 5000,
"1234567890_2_empty.mp4": 0,
"1234567890_4_tiny.jpg": 500,
}
r := Check(items, idx)
if r.Expected != 4 {
t.Errorf("Expected = %d, want 4", r.Expected)
}
if r.Present != 2 {
t.Errorf("Present = %d, want 2 (the non-empty ones)", r.Present)
}
if !slices.Equal(r.Absent, []int{3}) {
t.Errorf("Absent = %v, want [3]", r.Absent)
}
if !slices.Equal(r.ZeroByte, []int{2}) {
t.Errorf("ZeroByte = %v, want [2]", r.ZeroByte)
}
// Sub-1KiB is reported but still counted present: some real media is
// genuinely that small, so retrying it would loop forever.
if len(r.Tiny) != 1 || r.Tiny[0].MessageID != 4 {
t.Errorf("Tiny = %v, want just message 4", r.Tiny)
}
if slices.Contains(r.Todo(), 4) {
t.Error("a tiny file must not be queued for another fetch")
}
// Zero-byte files are retried: rclone overwrites a size-mismatched
// destination, so fetching again repairs them.
if want := []int{2, 3}; !slices.Equal(r.Todo(), want) {
t.Errorf("Todo() = %v, want %v", r.Todo(), want)
}
if r.Complete() {
t.Error("Complete() = true with work outstanding")
}
}
func TestCheckCompleteWhenEverythingIsPresent(t *testing.T) {
items := []tgsource.Item{item(1, "a.mp4", 10), item(2, "b.mp4", 20)}
idx := fakeIndex{"1234567890_1_a.mp4": 10, "1234567890_2_b.mp4": 20}
r := Check(items, idx)
if !r.Complete() {
t.Fatalf("Complete() = false, Todo() = %v", r.Todo())
}
var sb strings.Builder
r.Write(&sb)
if !strings.Contains(sb.String(), "COMPLETE") {
t.Errorf("report should say COMPLETE, got:\n%s", sb.String())
}
}
// The message that motivated the rewrite: the wanted name has a doubled '!',
// the remote holds the collapsed one that tdl's template wrote. It must be
// absent, and the stale copy must be surfaced rather than silently accepted.
func TestCheckReportsMisnamedCopiesAsAbsent(t *testing.T) {
items := []tgsource.Item{item(4242, "Pipe her!! And by her, we mean pipeperr! 1080p.mp4", 966444937)}
stale := "1234567890_4242_Pipe her! And by her, we mean pipeperr! 1080p.mp4"
idx := fakeIndex{stale: 966444937}
r := Check(items, idx)
if !slices.Equal(r.Absent, []int{4242}) {
t.Errorf("Absent = %v, want [4242] — a near-miss name is a different file", r.Absent)
}
if r.Present != 0 {
t.Errorf("Present = %d, want 0", r.Present)
}
if len(r.Misnamed) != 1 || !slices.Equal(r.Misnamed[0].Found, []string{stale}) {
t.Fatalf("Misnamed = %+v, want the stale copy reported", r.Misnamed)
}
var sb strings.Builder
r.Write(&sb)
out := sb.String()
for _, want := range []string{"stored under a different name", stale, "delete the stale copies"} {
if !strings.Contains(out, want) {
t.Errorf("report should mention %q, got:\n%s", want, out)
}
}
}
// A name that cannot be written to a path is never counted present and never
// silently skipped — it is reported and queued, so it stays visible.
func TestCheckFlagsUnsafeNames(t *testing.T) {
items := []tgsource.Item{item(7, "../../../.config/rclone/rclone.conf", 100)}
r := Check(items, fakeIndex{})
if len(r.Unsafe) != 1 || r.Unsafe[0].MessageID != 7 {
t.Fatalf("Unsafe = %+v, want message 7 flagged", r.Unsafe)
}
if !slices.Equal(r.Absent, []int{7}) {
t.Errorf("Absent = %v, want [7]", r.Absent)
}
if r.Present != 0 {
t.Errorf("Present = %d, want 0", r.Present)
}
var sb strings.Builder
r.Write(&sb)
if !strings.Contains(sb.String(), "unsafe filenames") {
t.Errorf("report should flag the unsafe name, got:\n%s", sb.String())
}
}
// An unsafe name must be rejected even when something is stored under that
// message id: the index is not the authority on whether a name is writable.
func TestCheckUnsafeNameIsNeverPresent(t *testing.T) {
items := []tgsource.Item{item(7, "sub/dir.mp4", 100)}
idx := fakeIndex{"1234567890_7_sub/dir.mp4": 100}
if r := Check(items, idx); r.Present != 0 || len(r.Unsafe) != 1 {
t.Errorf("Present = %d, Unsafe = %+v; an unsafe name must never count present", r.Present, r.Unsafe)
}
}
func TestReportTodoIsSorted(t *testing.T) {
r := Report{Absent: []int{4242, 3}, ZeroByte: []int{100}}
if want := []int{3, 100, 4242}; !slices.Equal(r.Todo(), want) {
t.Errorf("Todo() = %v, want %v", r.Todo(), want)
}
}