Files
telegram-exporter/internal/remote/fs.go
T
tiennm99 a81e7eaacb feat: download, upload and drive a chat to completion in one process
Phases 4 through 6: the two legs and the command that joins them.

Downloads go to <name>.part and are renamed only once complete, so a file
without the suffix is always whole. That is what lets the upload leg treat
"exists" as "finished" — the property run.sh could only approximate with a
filename convention plus an age guard, because it could not see inside tdl.

Every finished file is checked against the size Telegram reported, and that
check rather than the error is the authoritative signal. core's downloader logs
a failed transfer and returns nil, and its completion callback is deferred on
that named return, so a failure arrives indistinguishable from a success.
Trusting it would promote a truncated file and archive it as complete.

The disk cap is a semaphore over bytes. A download reserves its own size before
starting and releases it only after the upload confirms, so a slow remote
stalls downloads by itself. Blocking the iterator is safe because the
downloader calls it from its dispatch loop while workers run in a group, so a
blocked iterator never stops the uploads that free the space. Gone with it: the
du polling, the SIGSTOP and SIGCONT suspension, the min-age guard, the
temp-file filter and the sweep-failure counter.

A cap smaller than the largest file is refused up front. The semaphore could
never admit it, and a run blocked on a file it can never start looks exactly
like a stalled remote.

Uploads re-state each object to prove its size before the local copy is gone,
closing a gap where a truncated upload was only noticed by a later verify.

The destination is created before the chat is read. It is also the credentials
check, and doing it first means a bad destination fails in seconds rather than
after a full history walk.

One invocation converges: each item is checked against the index immediately
before download, so there are no passes and re-running is the resume path.
Options that no longer exist say what replaced them instead of failing as
unknown flags.

Verified end to end against the live chat and a scratch remote path: two files
downloaded, uploaded, confirmed present at the right size, staging left empty.
2026-09-06 19:21:46 +07:00

146 lines
5.1 KiB
Go

// Package remote owns the rclone side: resolving a destination and reporting on it.
package remote
import (
"context"
"errors"
"fmt"
"os"
"strings"
"sync"
"github.com/rclone/rclone/fs"
"github.com/rclone/rclone/fs/config"
"github.com/rclone/rclone/fs/config/configfile"
)
// Tunables are rclone settings this tool overrides.
//
// The values exist because of pikpak: it commits an upload as a server-side
// async task, and rclone abandons a still-pending one once its low-level retries
// run out, failing a transfer that would have succeeded. Fewer parallel
// transfers keep that queue short; more retries wait it out.
type Tunables struct {
Transfers int
LowLevelRetries int
}
// DefaultTunables are the values the shell pipeline settled on for pikpak.
func DefaultTunables() Tunables {
return Tunables{Transfers: 2, LowLevelRetries: 20}
}
// installOnce guards configfile.Install, which swaps unsynchronised package
// globals in rclone's config package. Calling it twice is harmless on its own,
// but racing it against an Fs resolution is not.
var installOnce sync.Once
// Init loads the user's rclone.conf and applies tunables to a derived context.
//
// Environment variables still win: rclone reads RCLONE_TRANSFERS and friends
// into its global config at package init, and fs.AddConfig copies that, so
// skipping the assignment when the variable is set preserves the operator's
// value.
//
// The config is loaded here, explicitly, because rclone's lazy path is fatal:
// config.LoadedData() calls fs.Fatalf on a config file it cannot parse or
// decrypt, and fs.Fatalf calls os.Exit(1) — past every defer, and with an exit
// code this tool defines as "incomplete", which would send a driver into an
// endless retry. Loading up front turns that into an ordinary error.
func Init(ctx context.Context, t Tunables) (context.Context, error) {
var err error
installOnce.Do(func() {
configfile.Install()
if lerr := config.Data().Load(); lerr != nil && !errors.Is(lerr, config.ErrorConfigFileNotFound) {
err = fmt.Errorf("cannot read rclone config %q: %w "+
"(an encrypted config needs RCLONE_CONFIG_PASS)", config.GetConfigPath(), lerr)
}
})
if err != nil {
return ctx, err
}
ctx, ci := fs.AddConfig(ctx)
if !envSet("RCLONE_TRANSFERS") && t.Transfers > 0 {
ci.Transfers = t.Transfers
}
if !envSet("RCLONE_LOW_LEVEL_RETRIES") && t.LowLevelRetries > 0 {
ci.LowLevelRetries = t.LowLevelRetries
}
return ctx, nil
}
// Resolve opens a destination given as an rclone REMOTE:PATH.
//
// rclone itself is the authority on what resolves: a remote can come from
// rclone.conf, from RCLONE_CONFIG_<NAME>_* environment variables with no config
// entry at all, from an inline `:type,opt=val:` connection string, or from a
// parameterised name like `pikpak,chunk_size=10M:path`. Pre-screening the name
// against the config sections would reject the last three, so the call is made
// first and the friendly "here is what you have configured" message is produced
// only for the one error that means the name was never defined.
func Resolve(ctx context.Context, remote string) (fs.Fs, error) {
if remote == "" {
return nil, fmt.Errorf("a destination remote is required (REMOTE:PATH)")
}
if !strings.Contains(remote, ":") {
return nil, fmt.Errorf("remote %q is not in rclone REMOTE:PATH form", remote)
}
f, err := fs.NewFs(ctx, remote)
if err != nil {
if errors.Is(err, fs.ErrorNotFoundInConfigFile) {
return nil, fmt.Errorf("rclone remote %q is not configured; configured remotes: %s",
remote, strings.Join(sections(), ", "))
}
return nil, fmt.Errorf("cannot reach %q — check credentials and connectivity: %w", remote, err)
}
return f, nil
}
// EnsureDir creates the destination if it is not there yet.
//
// A destination that does not exist yet is the normal case for a first run, and
// the shell pipeline created it up front for the same reason — the call doubles
// as the reachability and credentials check, since a remote that refuses a
// mkdir will refuse the uploads too. Doing it before the chat is read means a
// bad destination fails in seconds rather than after a full history walk.
func EnsureDir(ctx context.Context, f fs.Fs) error {
if err := f.Mkdir(ctx, ""); err != nil {
return fmt.Errorf("cannot create %q — check credentials and connectivity: %w", f.String(), err)
}
return nil
}
// FreeBytes reports free space on the remote.
//
// Backends without quota reporting return ok=false rather than an error: the
// shell pipeline treated an unanswerable quota as "unlimited" so a backend that
// cannot report never blocks a run, and that behaviour is preserved.
func FreeBytes(ctx context.Context, f fs.Fs) (free int64, ok bool) {
about := f.Features().About
if about == nil {
return 0, false
}
usage, err := about(ctx)
if err != nil || usage == nil || usage.Free == nil {
return 0, false
}
return *usage.Free, true
}
// sections lists configured remote names. Safe only after Init has loaded the
// config without error.
func sections() []string {
out := config.FileSections()
if len(out) == 0 {
return []string{"(none)"}
}
return out
}
func envSet(key string) bool {
_, ok := os.LookupEnv(key)
return ok
}