From bd816156ebc216a8b6ffa76d6c4a5c0bd8497379 Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Sun, 6 Sep 2026 19:37:30 +0700 Subject: [PATCH] fix: join download workers before closing the upload channel core's Download returns without waiting on its worker group when the iterator reports an error, so surfacing one through Iter.Err left workers sending into a channel the caller had already closed. The iterator now always reports a nil error and stashes the real one, read after Download returns. The circuit breaker cancelled only the upload context, which left downloads running full speed against a remote refusing them: every file stayed in staging and every reservation came back, so a broken remote filled local disk faster than a working one. It now stops the download iterator instead. A failed move leaves the local copy in place, so the byte reservation cannot be handed back until the file is removed. A confirmed-short object is deleted rather than left under a name verification would count as archived forever. Also: release the reservation when opening the staging file fails, report results even when an upload errored, bound the recorded errors, and drop the reporter's lock before writing so terminal latency cannot throttle downloads. --- README.md | 399 ++++++++++------------------- cmd/tgexport/sync.go | 38 +-- go.mod | 2 +- internal/pipeline/download.go | 32 ++- internal/pipeline/download_test.go | 123 ++++++++- internal/pipeline/elem.go | 69 +++-- internal/pipeline/pipeline.go | 94 ++++--- internal/pipeline/upload.go | 27 +- internal/report/progress.go | 13 +- 9 files changed, 439 insertions(+), 358 deletions(-) diff --git a/README.md b/README.md index b22174f..5cc42ae 100644 --- a/README.md +++ b/README.md @@ -1,305 +1,176 @@ # telegram-exporter -Export a Telegram chat's media to **any rclone remote** — S3, Google Drive, -Dropbox, Backblaze B2, SFTP, WebDAV, or anything else rclone supports — using -far less local disk than the chat's total size. +Archive a Telegram chat's media to **any rclone remote** — S3, Google Drive, +Dropbox, Backblaze B2, SFTP, WebDAV, pikpak, or anything else rclone supports — +using far less local disk than the chat's total size. -`run.sh` runs [tdl](https://github.com/iyear/tdl) and -[rclone](https://rclone.org/) as a rolling pipeline: tdl downloads into a small -staging directory while rclone concurrently moves finished files to the remote -and deletes the local copies. Local disk only ever holds the files in flight -plus one sync interval of throughput, so a multi-terabyte chat exports fine on a -small disk. Telegram caps a single file at 2 GB (4 GB from premium uploaders), -so a few dozen GB of staging covers the worst case regardless of chat size. +`tgexport` embeds [tdl](https://github.com/iyear/tdl) and +[rclone](https://rclone.org/) as libraries and runs both halves in one process. +Files are downloaded into a small staging directory and uploaded the moment each +one finishes, so local disk only ever holds what is in flight. A multi-terabyte +chat archives fine on a small disk. Telegram caps a single file at 2 GB (4 GB +from premium uploaders), so a few dozen GB of staging covers the worst case +regardless of chat size. -This is needed because **tdl can only write to a local directory** — it has no -rclone integration and no remote destination of any kind (`tdl dl -d` takes a -filesystem path; `tdl --storage` is its session database, not an output target). -tdl downloads over MTProto with a user account, so Bot API limits do not apply: -full history is readable and there is no 20 MB download cap. +This exists because **tdl can only write to a local directory** — it has no +remote destination of any kind (`tdl dl -d` takes a filesystem path; `tdl +--storage` is its session database, not an output target). tdl downloads over +MTProto with a user account, so Bot API limits do not apply: full history is +readable and there is no 20 MB download cap. ## Requirements -- **tdl** — -- **rclone** — -- **bash**. On Windows, run under WSL or Git Bash. +- **Go 1.25+** to build, or a prebuilt binary. +- **tdl** — only for `tdl login`. +- **rclone** — only to configure a remote. + +Neither tool is invoked at run time; `tgexport` reads the session and config +they write. ## Setup Both steps are one-time. ```bash -# 1) log in to Telegram with your user account (phone + code + 2FA) -tdl login - -# 2) configure the destination. The interactive wizard covers every backend: -rclone config - -# ...or create one non-interactively, e.g. -rclone config create gdrive drive -rclone config create b2 b2 account=KEY_ID key=APP_KEY -rclone config create dav webdav url=https://dav.example.com/remote.php/dav/files/you \ - vendor=other user=YOU pass=SECRET - -rclone listremotes # confirm the name you will pass to -r +tdl login # writes the Telegram session tgexport reads +rclone config # define the destination remote +go build -o tgexport ./cmd/tgexport +./tgexport doctor -r myremote:archive ``` -Any rclone remote form works, including on-the-fly connection strings -(`:webdav,url=https://...:/path`). +`doctor` proves both halves work before a long run: it prints the logged-in +account, resolves the destination, and reports free space. + +### Build variants + +| Build | Backends | Size | +|---|---|---| +| `go build ./cmd/tgexport` | every rclone backend | ~92 MB | +| `go build -tags slim ./cmd/tgexport` | pikpak only | ~49 MB | + +A backend that is not compiled in does not exist at run time, so use the default +build unless the destination will never change. ## Usage ```bash -./run.sh -r gdrive:telegram/media -c @mygroup +./tgexport sync -c CHAT -r REMOTE:PATH [options] ``` -That is the whole flow. It exports the chat's message metadata to -`export.json`, then downloads and uploads concurrently until finished. Progress -and warnings go to stderr; press Ctrl-C at any point and it stops cleanly. +`CHAT` accepts a numeric id as printed by `tdl chat ls`, a username with or +without `@`, or a `t.me`/`tg://` link. A Bot API `-100…` id is converted +automatically. A link to a single *message* is refused — it names a message, not +a chat. + +```bash +# archive a chat, capping staging at 40 GiB +./tgexport sync -c @mychannel -r gdrive:telegram/media -m 40G + +# check completeness without downloading anything +./tgexport verify -c @mychannel -r gdrive:telegram/media + +# list what the chat holds +./tgexport list -c @mychannel +``` ### Options -| Flag | Meaning | -|------|---------| -| `-r REMOTE:PATH` | **Required.** rclone destination, e.g. `gdrive:telegram/media`, `s3:bucket/tg`, `dav:tg-export` | -| `-c CHAT` | Chat to export when the JSON does not exist yet — id, username, or link (see below) | -| `-f FILE` | Export JSON to download from (default `export-.json` with `-c`, else `export.json`) | -| `-d DIR` | Staging directory (default `./staging`) | -| `-i SECONDS` | Seconds between rclone sweeps (default `60`) | -| `-a AGE` | rclone `--min-age`, a second guard against moving files still being written (default `45s`) | -| `-m SIZE` | Cap the staging directory at `SIZE` (`K`/`M`/`G`/`T`, binary), e.g. `40G`. Unset means no cap (see below) | -| `-h` | Help | +| Flag | Default | Meaning | +|---|---|---| +| `-c` | — | chat id, username, or link (required) | +| `-r` | — | rclone destination, `REMOTE:PATH` (required) | +| `-d` | `./staging` | staging directory for files in flight | +| `-m` | no cap | cap staging at a size, e.g. `40G` | +| `--threads` | 4 | connections per file | +| `--limit` | 2 | files downloading at once | +| `--uploads` | 2 | files uploading at once | +| `--min-free` | 5 | stop if the remote has fewer than this many GiB free | +| `--limit-items` | 0 | stop after N files; for smoke tests | +| `--confirm` | true | re-state each uploaded file to prove its size | +| `--takeout` | true | use a takeout session | +| `-n` | `default` | tdl session namespace | -### Identifying the chat +### Exit codes -`-c` accepts every form tdl understands, plus one it doesn't: +| Code | Meaning | +|---|---| +| 0 | complete | +| 1 | ran, but files remain | +| 2 | usage error | +| 3 | remote or Telegram failure | +| 130 / 143 | interrupted (SIGINT / SIGTERM) | -| Form | Example | -|------|---------| -| Numeric id, as printed by `tdl chat ls` | `-c 1697797156` | -| Username, with or without `@` | `-c @mygroup` / `-c mygroup` | -| Public link | `-c https://t.me/mygroup` / `-c t.me/mygroup` | -| Deep link | `-c 'tg://resolve?domain=mygroup'` | -| **Bot API id** (converted for you) | `-c -1001697797156` → `1697797156` | +## How it works -tdl resolves a numeric argument as an MTProto id and anything else through -gotd's resolver. MTProto has no `-100` prefix, so a Bot API id would otherwise -fail to resolve; the script strips it and logs the conversion. +Re-running is the resume path. Each item is checked against a listing of the +remote immediately before download, so an interrupted run picks up where it left +off and a completed one downloads nothing. -A **message** link is rejected — `-c` names a chat, not a message: +**Filenames.** Every file is stored as `{DialogID}_{MessageID}_{FileName}`, where +`FileName` is exactly what Telegram reports. One function derives that string, +and the same string is used both to ask whether the file is already archived and +to write it — so the two can never disagree. -``` -$ ./run.sh -r gdrive:tg -c https://t.me/mygroup/123 -error: -c takes a chat, not a message link — pass the chat's username or id -``` +That last point is the reason this program exists. Its predecessor derived the +name twice: `tdl chat export` wrote the raw name into a JSON, while `tdl dl` +rendered it through a template applying `filenamify`, which rewrites characters a +filesystem rejects and collapses runs of `!`. A file whose name contained `!!` +was looked up under one name and stored under another, so the verifier never +found it and re-fetched it on every pass — forever, at 966 MB a time. -Run `tdl chat ls` to see ids and usernames side by side. +Note the consequence: names are **not** run through `filenamify`, so they are not +byte-compatible with what the old shell pipeline wrote. A file it stored under a +rewritten name will not be recognised and gets fetched again. -Each chat gets its own export file by default (`export-mygroup.json`, -`export-1697797156.json`), so exporting a second chat from the same directory -never reuses the first one's JSON. When the file already exists it is reused and -the script says so — delete it to re-export. +**Disk.** `-m` is a byte budget. A download reserves its own size before starting +and releases it only once the upload is confirmed, so when the remote is slow the +downloads pause on their own. The cap must exceed the largest single file, and a +cap that does not is refused at startup rather than discovered as a hang. -Anything after `--` is passed straight to `tdl dl`: +**Integrity.** A download is written to `.part` and renamed only once its +size matches what Telegram reported, so a file without the suffix is always +whole. Uploads are re-stated afterwards to prove they arrived at the right size, +before the local copy is gone. -```bash -./run.sh -r gdrive:telegram/media -- -t 4 -l 1 # calmer parallelism, fewer flood waits -./run.sh -r gdrive:telegram/media -- -i mp4,mkv # only these file extensions -./run.sh -r gdrive:telegram/media -- -e jpg,png # skip these file extensions -``` +## Replacing the shell pipeline -tdl defaults to `-t 8 -l 4`, which is aggressive; lower it if you hit flood -waits on a large export. +Earlier versions of this repo were three bash scripts — `run.sh`, +`export-until-complete.sh` and `verify-export.sh` — driving `tdl` and `rclone` as +separate processes. Everything expensive in them existed to work around the fact +that neither process could see the other's state: a staging directory polled with +`du -sk`, an `--min-age` guard, a `*.tmp` exclusion, `SIGSTOP`/`SIGCONT` to +enforce the disk cap, a sweep-failure counter, and an outer loop that re-verified +and re-narrowed a JSON export between passes. -### Tuning the upload +One process needs none of it. Completion is a function returning; the cap is a +semaphore. Some hard-won details were worth keeping, and are: -rclone reads every one of its flags from an environment variable, so the upload -side is tunable without touching the script: +- **pikpak commits uploads as a server-side async task**, and rclone abandons a + still-pending one when its low-level retries run out. `transfers=2` and + `low-level-retries=20` are the defaults here for that reason. Environment + overrides still win. +- **A backend with no quota API is treated as unlimited**, so it never blocks a + run. +- **Zero-byte files count as missing** — rclone overwrites a size-mismatched + destination, so re-running repairs them — while files under 1 KiB are reported + but trusted, since some real media genuinely is that small. -```bash -RCLONE_TRANSFERS=8 RCLONE_BWLIMIT=20M ./run.sh -r s3:bucket/tg -c @mygroup -``` +Flags that disappeared are recognised and explain what replaced them: -`run.sh` sets two of those itself, and only when the caller has not: +| Old | Why it is gone | +|---|---| +| `-i` | no sweep interval; uploads start when a download finishes | +| `-a` | no `--min-age`; completion is observed, not inferred | +| `-f` | no export JSON; the chat is read live, so names cannot go stale | +| `-p` | no passes; one invocation converges | +| `-q` | renamed `--min-free` | -| Variable | Default here | rclone's own default | Why | -|----------|--------------|----------------------|-----| -| `RCLONE_TRANSFERS` | `2` | `4` | Backends that commit an upload as a server-side async task queue those tasks; less parallelism keeps the queue short | -| `RCLONE_LOW_LEVEL_RETRIES` | `20` | `10` | Each retry re-polls a pending task, so a slow commit is waited out instead of failing the transfer | +## Notes -Both exist because of one failure mode. On pikpak an upload finishes in two -phases — rclone sends the bytes, then a server-side task must reach -`PHASE_TYPE_COMPLETE`. rclone waits 500 ms and then polls, giving up after -`--low-level-retries` attempts with: - -``` -ERROR : : Failed to copy: can't verify the task is completed: ... Phase:"PHASE_TYPE_PENDING" -``` - -Nothing is lost when that happens — the message is followed by `Not deleting -source as copy failed`, the file stays in staging and the next sweep retries -it. But it wastes the upload, and it counts against a `-m` cap, since a file -that keeps failing can never be drained. Raise the retries further if you still -see it. - -### Capping the staging directory - -Without `-m`, staging grows whenever tdl downloads faster than rclone uploads, -which on a fast connection and a slow remote can mean tens of GB between -sweeps. `-m` puts a ceiling on it: - -```bash -./run.sh -r s3:bucket/tg -c @mygroup -m 40G -``` - -Staging size is checked every 10 seconds, independently of `-i`. When it -reaches the cap, tdl is suspended with `SIGSTOP` and rclone sweeps until -staging is back under it, then tdl is resumed — it reconnects on its own and -`--continue` picks its `.tmp` files back up. Because the checks are periodic, -the cap is a high-water mark rather than a hard limit: staging can overshoot by -up to ten seconds of download throughput before the gate closes. - -Only finished files can be drained, so the cap has to exceed what the -concurrent downloads hold — at most `-l` times 2 GB (4 GB from premium -uploaders). With the default `-l 2` anything from ~10 GB up is safe; below -that, the drain cannot clear the cap and the run logs a warning on every check -instead of throttling. - -### Exporting a subset - -Generate the JSON yourself when you want a narrower export, then point `-f` at -it: - -```bash -tdl chat export -c @mygroup -T id -i 1000,5000 --all --with-content -o part.json -./run.sh -r gdrive:telegram/media -f part.json -``` - -`tdl chat export` takes `-T time|id|last` with `-i` as the range, and `-f` as an -expression filter over message fields (`-f -` lists the available fields). - -## Sweep output - -The periodic sweeps are silent — they run every `-i` seconds alongside tdl's own -output, and narrating each one would drown it. The sweeps that run **once at the -end** do report progress, since they can move the whole staging directory with -nothing else on screen: - -- the exit sweep on Ctrl-C, `SIGTERM`, or a tdl failure (`sweeping completed - files before exit`); -- the final sweep after tdl finishes successfully. - -On a terminal that is rclone's redrawn `--progress` bar. When output is -redirected to a log it becomes a one-line stats summary every 30s -(`--stats 30s --stats-one-line --stats-log-level NOTICE`) — rclone logs stats at -INFO, so raising just the stats to NOTICE avoids the line-per-file spam that -`-v` would add. - -To show progress on every sweep instead, rclone reads its flags from the -environment: - -```bash -RCLONE_PROGRESS=true ./run.sh -r s3:bucket/tg -c @mygroup -``` - -## Resuming - -Re-run the same command. Both legs resume independently and nothing is -downloaded or uploaded twice. - -Keep the same `export.json` between runs: `--skip-same` compares against the -**staging** directory, which is empty once files have moved to the remote, so -cross-run deduplication rests on tdl's own `--continue` tracking. If you must -start from a fresh export, narrow it to the missing message-id range -(`-T id -i ,`) rather than re-downloading everything. - -## Verifying an export - -`run.sh` finishes when tdl finishes, which is not the same as every file having -arrived: a dropped session, a stalled remote, or an interrupted pass all leave -gaps. `verify-export.sh` settles it by rebuilding the filename tdl produces for -each media message in the export JSON and checking the remote for it. - -```bash -./verify-export.sh -f export-mygroup.json -r remote:telegram/media -``` - -``` -messages in export : 18193 - text-only (skip) : 38 - media expected : 12000 -present and intact : 12000 - absent : 0 - zero-byte : 0 - -COMPLETE: every media message is present and non-empty. -``` - -Messages with no media are skipped; they carry an empty `file` and were never -download targets. A zero-byte file counts as missing, because rclone overwrites a -size-mismatched destination and a retry repairs it. Files under 1 KiB are -reported but not retried, since some real media is genuinely that small. Exit -status is 0 when complete and 1 otherwise, with the outstanding message ids -written to `missing-ids.txt`. - -## Running until complete - -`export-until-complete.sh` drives `run.sh` in a loop: verify what is already -there, narrow the export to the ids still missing, run the pipeline on that -subset, and repeat. - -```bash -./export-until-complete.sh -r remote:telegram/media -c @mygroup -``` - -| Flag | Meaning | -|------|---------| -| `-r REMOTE:PATH` | **Required.** rclone destination | -| `-c CHAT` | Chat to export metadata for on the first pass | -| `-f FILE` | Export JSON (default `export-.json`) | -| `-d DIR` | Staging directory (default `./staging`) | -| `-i SECONDS` | rclone sweep interval (default `120`) | -| `-m SIZE` | Staging cap passed through to `run.sh`, e.g. `40G` | -| `-p N` | Maximum passes (default `30`) | -| `-q GIB` | Stop if remote free space falls below this (default `5`) | - -It stops when the verifier reports complete (exit `0`), when a pass fetches -nothing new (exit `1` — the remaining media is no longer available from -Telegram), when the remote runs low on space (exit `3`), or on Ctrl-C (exit -`130`, after the current pass shuts down cleanly). - -tdl's progress bar is shown when stdout is a terminal and suppressed when output -is redirected, so a log file stays readable without a flag. - -## What it guards against - -- **Partial uploads.** tdl writes `.tmp` and renames on completion, so - every sweep excludes `*.tmp`. Age alone is not a completion signal: a download - stalled by a flood wait stops touching its `.tmp`, which would then be - uploaded half-written and lose its resume point. -- **Directories vanishing under tdl.** `--delete-empty-src-dirs` runs only in - the final sweep, and the staging directory is recreated after every sweep. - Removing a directory under a running tdl makes it fail to create its next file. -- **Orphaned downloads.** tdl is stopped on exit, Ctrl-C, or `SIGTERM`, so no - download keeps running after the script is gone. -- **A failed run looking finished.** The unrestricted final sweep happens only - after tdl exits 0. An interrupted or crashed run gets the age-guarded sweep - and keeps staging for the next attempt. -- **A dead remote filling the disk.** Five consecutive rclone failures abort the - run instead of letting staging grow unbounded. -- **A fast connection filling the disk.** With `-m`, tdl is suspended whenever - staging reaches the cap and resumed once rclone has drained it, so download - throughput cannot outrun the upload leg. -- **Typos and bad credentials.** Before downloading anything, the remote must be - present in `rclone listremotes` (skipped for connection strings) and the - destination must be creatable, which proves both reachability and auth. - -Exit codes: `0` success, `2` usage error, `3` rclone failure, `130`/`143` -interrupted, anything else is tdl's own exit code. - -## Limits - -- Streaming with no staging at all (piping download chunks straight to the - remote) is not possible with tdl and would require custom code. -- The script is bash; the two tools it drives are cross-platform, but Windows - needs WSL or Git Bash. +- `tgexport` and the `tdl` CLI share one session store and cannot run against the + same namespace at once. Use `-n` for a second namespace if you need both. +- A partially downloaded file is not resumable across restarts — tdl's library + exposes no resume offset — so an interrupted run re-fetches whatever was in + flight, bounded by `--limit`. +- Everything is read-only against Telegram. Nothing is uploaded, deleted, or + marked read. diff --git a/cmd/tgexport/sync.go b/cmd/tgexport/sync.go index 207b139..d28b726 100644 --- a/cmd/tgexport/sync.go +++ b/cmd/tgexport/sync.go @@ -36,27 +36,27 @@ var retiredFlags = map[string]string{ // syncCmd archives a chat to a remote: read the chat, skip what is already // there, download and upload the rest, then report on the result. func syncCmd(ctx context.Context, args []string) error { - fs := flag.NewFlagSet("sync", flag.ContinueOnError) + flags := flag.NewFlagSet("sync", 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)") - staging = fs.String("d", "./staging", "staging directory for files in flight") - maxStaging = fs.String("m", "", "cap staging at this size, e.g. 40G (default: no cap)") - threads = fs.Int("threads", 4, "connections per file") - limit = fs.Int("limit", 2, "files downloading at once") - uploads = fs.Int("uploads", 2, "files uploading at once") - minFree = fs.Int64("min-free", 5, "stop if the remote has fewer than this many GiB free") - limitItems = fs.Int("limit-items", 0, "stop after this many files (0 means no limit)") - confirm = fs.Bool("confirm", true, "re-state each uploaded file to prove its size") - takeout = fs.Bool("takeout", true, "use a takeout session, as `tdl dl --takeout` did") - ns = fs.String("n", "default", "tdl session namespace") - dataDir = fs.String("storage", tdlkv.DefaultDir(), "tdl bolt storage directory") + chat = flags.String("c", "", "chat id, username, or t.me link (required)") + remoteArg = flags.String("r", "", "rclone destination, e.g. pikpak:archive (required)") + staging = flags.String("d", "./staging", "staging directory for files in flight") + maxStaging = flags.String("m", "", "cap staging at this size, e.g. 40G (default: no cap)") + threads = flags.Int("threads", 4, "connections per file") + limit = flags.Int("limit", 2, "files downloading at once") + uploads = flags.Int("uploads", 2, "files uploading at once") + minFree = flags.Int64("min-free", 5, "stop if the remote has fewer than this many GiB free") + limitItems = flags.Int("limit-items", 0, "stop after this many files (0 means no limit)") + confirm = flags.Bool("confirm", true, "re-state each uploaded file to prove its size") + takeout = flags.Bool("takeout", true, "use a takeout session, as `tdl dl --takeout` did") + ns = flags.String("n", "default", "tdl session namespace") + dataDir = flags.String("storage", tdlkv.DefaultDir(), "tdl bolt storage directory") ) for name, replacement := range retiredFlags { - fs.Var(retiredFlag{name, replacement}, name, "retired") + flags.Var(retiredFlag{name, replacement}, name, "retired") } - if err := fs.Parse(args); err != nil { + if err := flags.Parse(args); err != nil { if errors.Is(err, flag.ErrHelp) { return err } @@ -169,7 +169,11 @@ func syncCmd(ctx context.Context, args []string) error { fmt.Fprintf(os.Stderr, " message %d failed: %v\n", f.Item.MessageID, f.Err) } if runErr != nil { - return runErr + // Reported, not returned yet: a run that archived thousands of files + // and hit one transient upload error has still made progress, and + // suppressing the report would leave the operator — and any driver + // reading the exit code — unable to tell that from a total failure. + fmt.Fprintf(os.Stderr, "run ended early: %v\n", runErr) } // The remote is re-indexed rather than assumed: the run's own view of diff --git a/go.mod b/go.mod index c4cae8f..ee2e563 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( github.com/iyear/tdl/core v0.20.4 github.com/rclone/rclone v1.75.1 go.etcd.io/bbolt v1.5.0 + golang.org/x/sync v0.22.0 ) require ( @@ -227,7 +228,6 @@ require ( golang.org/x/mod v0.38.0 // indirect golang.org/x/net v0.58.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect - golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/term v0.45.0 // indirect golang.org/x/text v0.41.0 // indirect diff --git a/internal/pipeline/download.go b/internal/pipeline/download.go index e160c0a..5e4a0ba 100644 --- a/internal/pipeline/download.go +++ b/internal/pipeline/download.go @@ -8,6 +8,7 @@ import ( "os" "path/filepath" "strings" + "sync/atomic" "github.com/iyear/tdl/core/dcpool" "github.com/iyear/tdl/core/downloader" @@ -29,8 +30,14 @@ type DownloadOptions struct { Report func(Stats) // acquire reserves staging space before a download starts, blocking until - // there is room. Unset means no bound. + // there is room. Unset means no bound. release hands a reservation back for + // an item that never reaches a download. acquire func(context.Context, int64) error + release func(int64) + + // stop, when set, ends iteration cleanly from another goroutine — used to + // halt downloads once the destination has stopped accepting uploads. + stop *atomic.Bool // onReady hands a completed file to the upload leg; onFailed says nothing // was staged, so whatever acquire reserved must be given back. onReady func(tgsource.Item) @@ -39,12 +46,13 @@ type DownloadOptions struct { // Download fetches every item in seq into the staging directory. // -// Each file is written to .part and renamed to only once the -// downloader reports it complete, so a name without the suffix is always a -// whole file. That is what lets the upload half treat "the file exists" as -// "the file is finished" — the property the shell pipeline had to approximate -// with a filename convention plus an age guard, because it could not see -// inside tdl. +// Each file is written to .part and renamed to only once its size +// matches what Telegram reported, so a name without the suffix is always a whole +// file. Uploads are driven by completion rather than by scanning for that, but +// the invariant still matters: it is what makes a leftover file from an +// interrupted run safe to keep and a leftover .part safe to delete. The shell +// pipeline could only approximate it with a filename convention plus an age +// guard, because it could not see inside tdl. // // A failed item does not abort the run: it is recorded in the returned outcomes // and the rest continue, matching what a partial `tdl dl` pass did. @@ -61,6 +69,10 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo it := newElemIter(seq, o.Staging, o.Takeout) it.acquire = o.acquire + it.release = o.release + if o.stop != nil { + it.stopped = o.stop + } defer func() { _ = it.Close() }() prog := newProgress(func(e *elem, err error) error { @@ -85,6 +97,12 @@ func Download(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Downlo }).Download(ctx, o.Limit) outcomes, stats := prog.results() + // The iterator's failure is read only now, after Download has joined every + // worker. Reporting it through Iter.Err would have made Download skip that + // join entirely. + if err == nil { + err = it.failure + } return outcomes, stats, err } diff --git a/internal/pipeline/download_test.go b/internal/pipeline/download_test.go index 7940c46..9952a5b 100644 --- a/internal/pipeline/download_test.go +++ b/internal/pipeline/download_test.go @@ -1,6 +1,7 @@ package pipeline import ( + "context" "errors" "os" "path/filepath" @@ -141,12 +142,11 @@ func TestElemIterRejectsUnsafeNames(t *testing.T) { if it.Next(t.Context()) { t.Fatal("iterator accepted a name that escapes the staging directory") } - err := it.Err() - if err == nil { - t.Fatal("Err() = nil after rejecting an unsafe name") + if it.failure == nil { + t.Fatal("no failure recorded after rejecting an unsafe name") } - if !strings.Contains(err.Error(), "message 7") { - t.Errorf("error should name the message, got: %v", err) + if !strings.Contains(it.failure.Error(), "message 7") { + t.Errorf("failure should name the message, got: %v", it.failure) } } @@ -183,8 +183,8 @@ func TestElemIterOpensPartFilesAndPropagatesWalkErrors(t *testing.T) { if it.Next(t.Context()) { t.Fatal("Next() = true after a walk error") } - if !errors.Is(it.Err(), want) { - t.Errorf("Err() = %v, want %v", it.Err(), want) + if !errors.Is(it.failure, want) { + t.Errorf("failure = %v, want %v", it.failure, want) } }) } @@ -288,3 +288,112 @@ func TestFinishAcceptsExactSize(t *testing.T) { t.Errorf("final size = %d, want 2048", info.Size()) } } + +// Err must always report nil, however badly iteration went. +// +// core's Download skips wg.Wait entirely when Iter.Err is non-nil +// (downloader.go:65-68), returning while its workers are still running. The +// pipeline closes its upload channel as soon as Download returns, so a non-nil +// Err here means workers send on a closed channel and the process panics — +// on every Ctrl-C, since cancellation is one of the ways iteration stops. +func TestElemIterNeverReportsErrToTheDownloader(t *testing.T) { + staging := t.TempDir() + + cases := map[string]func() *elemIter{ + "walk error": func() *elemIter { + seq := func(yield func(tgsource.Item, error) bool) { + yield(tgsource.Item{}, errors.New("boom")) + } + return newElemIter(seq, staging, false) + }, + "unsafe name": func() *elemIter { + seq := func(yield func(tgsource.Item, error) bool) { + yield(testItem(t, 1, "../escape", 10), nil) + } + return newElemIter(seq, staging, false) + }, + "cancelled": func() *elemIter { + seq := func(yield func(tgsource.Item, error) bool) { + yield(testItem(t, 2, "a.mp4", 10), nil) + } + return newElemIter(seq, staging, false) + }, + } + + for name, build := range cases { + t.Run(name, func(t *testing.T) { + it := build() + defer func() { _ = it.Close() }() + + ctx := t.Context() + if name == "cancelled" { + cancelled, cancel := context.WithCancel(ctx) + cancel() + ctx = cancelled + } + + for it.Next(ctx) { + } + if err := it.Err(); err != nil { + t.Errorf("Err() = %v, want nil — a non-nil Err makes Download abandon its workers", err) + } + if it.failure == nil { + t.Error("the real failure was not stashed") + } + }) + } +} + +// A caller-set stop flag ends iteration without looking like a failure, which is +// how the circuit breaker halts downloads. +func TestElemIterStopsOnFlagWithoutRecordingFailure(t *testing.T) { + staging := t.TempDir() + seq := func(yield func(tgsource.Item, error) bool) { + for i := 1; i <= 5; i++ { + if !yield(testItem(t, i, "a.mp4", 10), nil) { + return + } + } + } + it := newElemIter(seq, staging, false) + defer func() { _ = it.Close() }() + + if !it.Next(t.Context()) { + t.Fatal("first Next() = false") + } + it.stopped.Store(true) + + if it.Next(t.Context()) { + t.Error("Next() = true after the stop flag was set") + } + if it.failure != nil { + t.Errorf("failure = %v, want nil — stopping is not a failure", it.failure) + } +} + +// A reservation must come back when the item never reaches a download, or the +// budget shrinks by that much for the rest of the run. +func TestElemIterReturnsReservationWhenOpenFails(t *testing.T) { + // A staging path that is a file, not a directory, makes OpenFile fail. + staging := filepath.Join(t.TempDir(), "not-a-dir") + if err := os.WriteFile(staging, []byte("x"), 0o600); err != nil { + t.Fatalf("seed: %v", err) + } + + seq := func(yield func(tgsource.Item, error) bool) { + yield(testItem(t, 1, "a.mp4", 4096), nil) + } + it := newElemIter(seq, staging, false) + defer func() { _ = it.Close() }() + + var acquired, released int64 + it.acquire = func(_ context.Context, n int64) error { acquired += n; return nil } + it.release = func(n int64) { released += n } + + if it.Next(t.Context()) { + t.Fatal("Next() succeeded with an unusable staging directory") + } + if acquired != released { + t.Errorf("acquired %d bytes but released %d — the reservation leaked", acquired, released) + } +} diff --git a/internal/pipeline/elem.go b/internal/pipeline/elem.go index d3c3464..703abe8 100644 --- a/internal/pipeline/elem.go +++ b/internal/pipeline/elem.go @@ -8,6 +8,7 @@ import ( "iter" "os" "path/filepath" + "sync/atomic" "github.com/gotd/td/tg" @@ -59,37 +60,57 @@ func finalPath(staging string, it tgsource.Item) string { // 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 + next func() (tgsource.Item, error, bool) + stopPull func() + staging string + takeout bool // acquire reserves staging space for the next item. Blocking here is what // makes backpressure work: core's Download calls Next from its dispatch - // loop, so a blocked Next stops new downloads starting without stopping the - // uploads that free the space. + // loop (downloader.go:38), so a blocked Next stops new downloads starting + // without stopping the uploads that free the space. acquire func(context.Context, int64) error + // release hands a reservation back when the item never reaches a download. + release func(int64) 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. + // failure holds why iteration stopped, and Err deliberately does not return + // it. core's Download skips wg.Wait entirely when Iter.Err is non-nil + // (downloader.go:65-68), abandoning workers that are still running — which + // would let this package tear down its upload channel underneath them. So + // Next reports "no more items" and the caller reads failure() afterwards, + // guaranteeing every worker has finished first. + failure error + // stopped ends iteration without an error, for a caller that has decided the + // run cannot usefully continue. Supplied by the caller so it can be set from + // another goroutine without racing on the iterator itself. + stopped *atomic.Bool + + // opened records every file handle. finish closes each one on the normal + // path, so this is not what keeps descriptors from leaking; it is the + // backstop for items that were opened but never reached finish, which is + // what an aborted iteration leaves behind. 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} + next, stopPull := iter.Pull2(seq) + return &elemIter{ + next: next, + stopPull: stopPull, + staging: staging, + takeout: takeout, + stopped: new(atomic.Bool), + } } func (i *elemIter) Next(ctx context.Context) bool { - if i.err != nil { + if i.failure != nil || i.stopped.Load() { return false } if err := ctx.Err(); err != nil { - i.err = err + i.failure = err return false } @@ -98,7 +119,7 @@ func (i *elemIter) Next(ctx context.Context) bool { return false } if err != nil { - i.err = err + i.failure = err return false } @@ -106,20 +127,25 @@ func (i *elemIter) Next(ctx context.Context) bool { // 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) + i.failure = fmt.Errorf("message %d: %w", item.MessageID, err) return false } if i.acquire != nil { if err := i.acquire(ctx, item.Size()); err != nil { - i.err = err + i.failure = 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) + // The reservation is handed back here because this item will never + // reach a download, so no OnDone will ever release it for us. + if i.release != nil { + i.release(item.Size()) + } + i.failure = fmt.Errorf("open destination for message %d: %w", item.MessageID, err) return false } i.opened = append(i.opened, f) @@ -129,11 +155,14 @@ func (i *elemIter) Next(ctx context.Context) bool { } func (i *elemIter) Value() downloader.Elem { return i.current } -func (i *elemIter) Err() error { return i.err } + +// Err always reports nil so core's Download reaches wg.Wait and joins its +// workers. See the failure field. +func (i *elemIter) Err() error { return nil } // Close releases the pull iterator and every file the walk opened. func (i *elemIter) Close() error { - i.stop() + i.stopPull() var firstErr error for _, f := range i.opened { if err := f.Close(); err != nil && firstErr == nil { diff --git a/internal/pipeline/pipeline.go b/internal/pipeline/pipeline.go index 2bcd1b4..af9b391 100644 --- a/internal/pipeline/pipeline.go +++ b/internal/pipeline/pipeline.go @@ -5,7 +5,10 @@ import ( "errors" "fmt" "iter" + "os" + "path/filepath" "sync" + "sync/atomic" "github.com/rclone/rclone/fs" "golang.org/x/sync/semaphore" @@ -53,6 +56,11 @@ func (r Result) Failed() []Outcome { return out } +// maxRecordedErrors bounds what a run keeps from a failing remote. Past this, +// the pattern is established and joining thousands of identical strings just +// makes the final message unreadable. +const maxRecordedErrors = 10 + // Run downloads every item and uploads each one as it completes. // // This is the whole reason for the rewrite. run.sh could not see inside tdl, so @@ -60,13 +68,17 @@ func (r Result) Failed() []Outcome { // `du -sk` every ten seconds, and enforced its disk cap by sending SIGSTOP and // SIGCONT to the tdl process. None of that exists here. Completion is a function // returning. The cap is a semaphore: a download acquires its own size before -// starting and releases it only once the upload has confirmed, so when the -// remote is slow the acquire blocks and downloads pause on their own. +// starting and releases it once the file is off local disk, so when the remote +// is slow the acquire blocks and downloads pause on their own. // -// Blocking in the iterator is safe by construction — core's Download calls -// Iter.Next from its dispatch loop while workers run in an errgroup, so a -// blocked Next stalls new work without stopping the uploads that free the budget -// (downloader.go:36-63). +// Blocking in the iterator is safe, but not for the reason it first appears. +// core's Download calls Iter.Next from its dispatch loop while workers run in an +// errgroup, so a blocked Next stalls new work without stopping the uploads that +// free the budget. What is *not* safe is reporting an error through Iter.Err: +// Download then returns without joining its workers (downloader.go:65-68), and +// tearing down the upload channel underneath them panics. So elemIter always +// reports a nil Err and stashes the real one, which Download's return +// guarantees is safe to read. func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (Result, error) { if o.Uploads <= 0 { o.Uploads = 1 @@ -84,37 +96,51 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R budget := newBudget(o.Budget) uploads := make(chan tgsource.Item, o.Uploads) - // Upload workers own the release side of the budget, so every path out of - // one — success, failure, cancellation — must release, or the run deadlocks - // with downloads waiting on space that is never freed. var ( - wg sync.WaitGroup - mu sync.Mutex - uploadErrs []error - streak int - tripped bool + wg sync.WaitGroup + mu sync.Mutex + errs []error + nErrs int + streak int + tripped bool ) - upCtx, tripRun := context.WithCancel(ctx) - defer tripRun() + + // stopDownloads ends the download side once the destination has stopped + // accepting work. It stops the iterator rather than cancelling a context, + // because cancelling only the uploads would leave downloads running at full + // speed against a remote that is refusing them — every file staying on disk, + // every reservation released on the way out. A broken remote would fill the + // local disk faster than a working one does. + stopDownloads := new(atomic.Bool) for range o.Uploads { wg.Add(1) go func() { defer wg.Done() for it := range uploads { - err := up.upload(upCtx, it) + err := up.upload(ctx, it) + + if err != nil { + // MoveFile leaves the local copy in place when it fails, so + // the reservation cannot simply be handed back — the bytes + // are still on disk. Removing the file first is what keeps + // the cap honest. + if rerr := os.Remove(filepath.Join(o.Staging, it.Name)); rerr != nil && !os.IsNotExist(rerr) { + err = errors.Join(err, fmt.Errorf("and it is still in staging: %w", rerr)) + } + } budget.release(it.Size()) mu.Lock() if err != nil { - uploadErrs = append(uploadErrs, err) + if nErrs < maxRecordedErrors { + errs = append(errs, err) + } + nErrs++ streak++ if streak >= o.MaxFailures && !tripped { - // A remote that fails this many times running is not - // going to recover on its own, and continuing just fills - // staging until the disk does. tripped = true - tripRun() + stopDownloads.Store(true) } } else { streak = 0 @@ -132,20 +158,24 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R Takeout: o.Takeout, Report: o.Report, acquire: budget.acquire, + release: budget.release, onReady: func(it tgsource.Item) { uploads <- it }, onFailed: func(it tgsource.Item) { // Nothing was staged, so the reservation has to come back here // instead of from an upload that will never happen. budget.release(it.Size()) }, + stop: stopDownloads, }) + // Safe only because Download joined its workers, which is guaranteed by + // elemIter.Err always being nil. close(uploads) wg.Wait() mu.Lock() - errs := append([]error(nil), uploadErrs...) - trip := tripped + joined := errors.Join(errs...) + trip, total := tripped, nErrs mu.Unlock() res := Result{Stats: stats, Outcomes: dlOutcomes} @@ -153,10 +183,10 @@ func Run(ctx context.Context, seq iter.Seq2[tgsource.Item, error], o Options) (R case dlErr != nil: return res, dlErr case trip: - return res, fmt.Errorf("stopping after %d consecutive upload failures: %w", - o.MaxFailures, errors.Join(errs...)) - case len(errs) > 0: - return res, errors.Join(errs...) + return res, fmt.Errorf("stopped after %d consecutive upload failures (%d total): %w", + o.MaxFailures, total, joined) + case total > 0: + return res, fmt.Errorf("%d upload(s) failed: %w", total, joined) } return res, nil } @@ -176,10 +206,10 @@ func (b *budget) acquire(ctx context.Context, n int64) error { if b.sem == nil { return nil } - // An item larger than the whole budget could never be admitted and would - // block forever, so it is refused with an error that says what to change. - // Callers validate up front too; this is the guard for an item whose size - // was not known then. + // An item larger than the budget cannot be admitted, and semaphore.Acquire + // handles that by blocking until the context is cancelled rather than + // failing — so there is no error to surface and no guard to add here. The + // real protection is validateBudget refusing such a run before it starts. if err := b.sem.Acquire(ctx, n); err != nil { return fmt.Errorf("waiting for %d bytes of staging space: %w", n, err) } diff --git a/internal/pipeline/upload.go b/internal/pipeline/upload.go index 89ed3eb..f38f108 100644 --- a/internal/pipeline/upload.go +++ b/internal/pipeline/upload.go @@ -3,6 +3,7 @@ package pipeline import ( "context" "fmt" + "time" "github.com/rclone/rclone/fs" "github.com/rclone/rclone/fs/operations" @@ -21,13 +22,15 @@ type uploader struct { // arrived at the expected size. // // MoveFile removes the local copy as part of the move, so a successful return -// means the file is on the remote and off local disk — which is what lets the -// byte budget be released. +// means the file is on the remote and off local disk. // -// Confirmation closes a gap the shell pipeline left open: there, a truncated -// upload was only noticed by a later verify pass, after the local copy was -// already gone. Re-stating the object costs one round trip per file and turns a -// silent corruption into a retry. +// A short object is deleted rather than left in place, and that is the part that +// matters. Verification matches on name and non-zero size, so a truncated object +// under the right name would be counted archived by this run and by every run +// after it — permanently, with the local copy already gone. Removing it turns a +// silent corruption into an absent file the next run fetches again. The shell +// pipeline had this hole too: it noticed a bad upload only at the next verify, +// by which point the evidence was the same. func (u *uploader) upload(ctx context.Context, it tgsource.Item) error { if err := operations.MoveFile(ctx, u.dst, u.local, it.Name, it.Name); err != nil { return fmt.Errorf("move %q to %s: %w", it.Name, u.dst.String(), err) @@ -41,7 +44,17 @@ func (u *uploader) upload(ctx context.Context, it tgsource.Item) error { return fmt.Errorf("confirm %q: %w", it.Name, err) } if got := obj.Size(); got != it.Size() { - return fmt.Errorf("confirm %q: remote has %d bytes, expected %d", it.Name, got, it.Size()) + err := fmt.Errorf("confirm %q: remote has %d bytes, expected %d", it.Name, got, it.Size()) + // Deleted on a fresh context: the run may already be shutting down, and + // leaving a plausible-looking short object behind is worse than the + // error that got us here. + delCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if derr := operations.DeleteFile(delCtx, obj); derr != nil { + return fmt.Errorf("%w (and it could not be removed: %v — delete it by hand "+ + "or verify will count it archived)", err, derr) + } + return err } return nil } diff --git a/internal/report/progress.go b/internal/report/progress.go index e458fa8..69cf2c7 100644 --- a/internal/report/progress.go +++ b/internal/report/progress.go @@ -41,10 +41,17 @@ func New(w io.Writer, total int, totalBytes int64) *Reporter { return &Reporter{w: w, tty: isTerminal(w), total: total, totalBytes: totalBytes, started: time.Now()} } -// Update renders a snapshot. Safe to call from several goroutines, and cheap -// enough to call on every progress callback. +// Update renders a snapshot. Safe to call from several goroutines. +// +// A contended update is dropped rather than queued. Every download worker calls +// this on each progress callback, so holding the lock across the write would +// make terminal latency — an ssh session with a slow link, say — throttle the +// downloads themselves. A skipped frame costs nothing; the next callback is +// milliseconds away and Finish always prints. func (r *Reporter) Update(s pipeline.Stats) { - r.mu.Lock() + if !r.mu.TryLock() { + return + } defer r.mu.Unlock() now := time.Now()