diff --git a/cmd/tgexport/main.go b/cmd/tgexport/main.go index eab509c..0a515fe 100644 --- a/cmd/tgexport/main.go +++ b/cmd/tgexport/main.go @@ -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 idsizename + verify Report whether a chat is fully archived on a remote Run 'tgexport -h' for command options. `) diff --git a/cmd/tgexport/verify.go b/cmd/tgexport/verify.go new file mode 100644 index 0000000..d79d1e7 --- /dev/null +++ b/cmd/tgexport/verify.go @@ -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...) +} diff --git a/internal/naming/naming_test.go b/internal/naming/naming_test.go index a96f8e4..a14d86c 100644 --- a/internal/naming/naming_test.go +++ b/internal/naming/naming_test.go @@ -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} } diff --git a/internal/naming/safe.go b/internal/naming/safe.go new file mode 100644 index 0000000..feebf6f --- /dev/null +++ b/internal/naming/safe.go @@ -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 +} diff --git a/internal/naming/safe_test.go b/internal/naming/safe_test.go new file mode 100644 index 0000000..531202d --- /dev/null +++ b/internal/naming/safe_test.go @@ -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))) + } +} diff --git a/internal/pipeline/elem.go b/internal/pipeline/elem.go new file mode 100644 index 0000000..30e39fa --- /dev/null +++ b/internal/pipeline/elem.go @@ -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 +} diff --git a/internal/remote/index.go b/internal/remote/index.go new file mode 100644 index 0000000..73bd950 --- /dev/null +++ b/internal/remote/index.go @@ -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 +} diff --git a/internal/remote/index_test.go b/internal/remote/index_test.go new file mode 100644 index 0000000..5abd8c4 --- /dev/null +++ b/internal/remote/index_test.go @@ -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) + } +} diff --git a/internal/verify/verify.go b/internal/verify/verify.go new file mode 100644 index 0000000..c72c72e --- /dev/null +++ b/internal/verify/verify.go @@ -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") +} diff --git a/internal/verify/verify_test.go b/internal/verify/verify_test.go new file mode 100644 index 0000000..6978b77 --- /dev/null +++ b/internal/verify/verify_test.go @@ -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) + } +}