Dead, vestigial, and over-built code with no behavior change (except the
bool-vocabulary trim, a deliberate behavior tightening):
- httpdl: drop the vestigial sort+copy in validSegs (work-stealing was
removed in 270812d, so segments always tile [0,total) ascending now);
validate order in place, which also rejects a scrambled sidecar. Inline
the endOff() accessor to s.end.
- cli: shrink the boolean vocabulary to true/false, which is all that's
accepted. The invented yes/no/1/0/on/off spellings are gone, so e.g.
--enable-dht=yes (or yes in got.conf) is now a usage error.
- download: remove the dead, never-called Stat.Done() method.
- progress: drop speed(), a pure alias of humanSize; call humanSize directly.
- main: replace the jobsHolder mutex box with atomic.Pointer[[]job] for the
forced-exit session save; the capability stays, the lock and type go.
- bt: collapse isAddrInUse's X||X (missinggo.IsAddrInUse is that exact
string match); drop the now-unused missinggo import (tidy -> indirect).
- option: note that prealloc and falloc are identical here.
550 lines
18 KiB
Go
550 lines
18 KiB
Go
// Package bt downloads a torrent or magnet link. The peer wire protocol, DHT,
|
|
// trackers, magnet metadata, piece selection and choking are all carried by
|
|
// github.com/anacrolix/torrent; this package's job is to present that engine as
|
|
// a download.Download. Within a single run, one shared torrent.Client serves
|
|
// every download, so they share one listen port, DHT node and piece-store
|
|
// rather than fighting over them.
|
|
package bt
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
alog "github.com/anacrolix/log"
|
|
"github.com/anacrolix/torrent"
|
|
"github.com/anacrolix/torrent/metainfo"
|
|
"github.com/hanbok/got/download"
|
|
"golang.org/x/time/rate"
|
|
)
|
|
|
|
// ClientConfig holds the run-wide settings used to build the shared client.
|
|
type ClientConfig struct {
|
|
Dir string
|
|
ListenPortSpec string // raw --listen-port spec ("6881-6999,7000"); empty lets the OS choose.
|
|
DHT bool
|
|
MaxPeers int // max peer connections per torrent (0 = library default)
|
|
OverallDown *rate.Limiter // global down limit (may be nil)
|
|
OverallUp *rate.Limiter // global up limit (may be nil)
|
|
DownLimit int64 // per-run bytes/s, 0 = unlimited (used only without a global limiter)
|
|
UpLimit int64
|
|
UserAgent string
|
|
DisableIPv6 bool // force IPv4-only (--disable-ipv6)
|
|
}
|
|
|
|
// ParsePorts expands a --listen-port spec into a flat list of candidate ports.
|
|
// The spec is a comma-separated list of single ports and a-b ranges
|
|
// (e.g. "6881-6889,6999"). Out-of-range and unparsable entries are
|
|
// skipped; an empty or all-bad spec yields nil.
|
|
func ParsePorts(spec string) []int {
|
|
var ports []int
|
|
for _, part := range strings.Split(spec, ",") {
|
|
part = strings.TrimSpace(part)
|
|
if part == "" {
|
|
continue
|
|
}
|
|
lo, hi, isRange := strings.Cut(part, "-")
|
|
a, err := strconv.Atoi(strings.TrimSpace(lo))
|
|
if err != nil {
|
|
// Unparsable low bound: drop the whole entry, symmetrically with the
|
|
// high bound below (an out-of-range a is still filtered per-port).
|
|
continue
|
|
}
|
|
if !isRange {
|
|
if a >= 1 && a <= 65535 {
|
|
ports = append(ports, a)
|
|
}
|
|
continue
|
|
}
|
|
b, err := strconv.Atoi(strings.TrimSpace(hi))
|
|
if err != nil {
|
|
continue
|
|
}
|
|
// Clamp to the valid port space before iterating: a parseable but oversized
|
|
// spec like "1-9999999999" would otherwise spin for ~10^10 iterations and
|
|
// hang startup, even though only [1,65535] can ever be appended.
|
|
if a < 1 {
|
|
a = 1
|
|
}
|
|
if b > 65535 {
|
|
b = 65535
|
|
}
|
|
for i := a; i <= b; i++ {
|
|
ports = append(ports, i)
|
|
}
|
|
}
|
|
return ports
|
|
}
|
|
|
|
// quietSlogger filters anacrolix's structured logging down to Error and above,
|
|
// so the routine webseed/peer/tracker warning chatter (the 403/429 noise) no
|
|
// longer interleaves with — and corrupts — the live progress display. Genuine
|
|
// errors still reach w.
|
|
func quietSlogger(w io.Writer) *slog.Logger {
|
|
return slog.New(slog.NewTextHandler(w, &slog.HandlerOptions{Level: slog.LevelError}))
|
|
}
|
|
|
|
// NewClient builds the shared anacrolix client from the run-wide config. When a
|
|
// listen-port spec is given it tries each candidate port in turn, advancing past
|
|
// any that is already in use, and falls back to an OS-chosen port if none bind,
|
|
// so a single busy port does not disable BitTorrent entirely.
|
|
//
|
|
// Two per-torrent features have no equivalent in anacrolix/torrent v1.61.0,
|
|
// which exposes only client-wide knobs; both must therefore be configured run-wide
|
|
// here rather than attached to an individual Torrent:
|
|
//
|
|
// 1. No per-Torrent rate limiter. Throttles are read from the client-wide
|
|
// DownloadRateLimiter/UploadRateLimiter (no TorrentSpec rate field exists), so
|
|
// per-torrent --max-download-limit / --max-upload-limit cannot be
|
|
// honored individually; the effective cap is the shared limiter set below.
|
|
// 2. No per-Torrent DHT/PEX/LSD toggle for BEP27 private torrents. Those are
|
|
// client-wide only (NoDHT, DisablePEX, etc.), so a private torrent's leak
|
|
// prevention must come from a run-wide setting (e.g. --enable-dht=false).
|
|
func NewClient(cfg ClientConfig) (*torrent.Client, error) {
|
|
c := torrent.NewDefaultClientConfig()
|
|
// Keep the library's own logging out of the readout. Webseed/peer warnings go
|
|
// through slog, the rest through the legacy analog logger, so filter both to
|
|
// Error and above (the analog default is Warning, which is what leaks today).
|
|
c.Slogger = quietSlogger(os.Stderr)
|
|
c.Logger = alog.Default.FilterLevel(alog.Error)
|
|
c.DataDir = cfg.Dir
|
|
if c.DataDir == "" {
|
|
c.DataDir = "."
|
|
}
|
|
c.Seed = true
|
|
c.NoDHT = !cfg.DHT
|
|
if cfg.MaxPeers > 0 {
|
|
c.EstablishedConnsPerTorrent = cfg.MaxPeers
|
|
}
|
|
if cfg.UserAgent != "" {
|
|
c.HTTPUserAgent = cfg.UserAgent
|
|
}
|
|
c.DisableIPv6 = cfg.DisableIPv6
|
|
switch {
|
|
case cfg.OverallDown != nil:
|
|
c.DownloadRateLimiter = cfg.OverallDown
|
|
case cfg.DownLimit > 0:
|
|
c.DownloadRateLimiter = rate.NewLimiter(rate.Limit(cfg.DownLimit), download.LimiterBurst(cfg.DownLimit))
|
|
}
|
|
switch {
|
|
case cfg.OverallUp != nil:
|
|
c.UploadRateLimiter = cfg.OverallUp
|
|
case cfg.UpLimit > 0:
|
|
c.UploadRateLimiter = rate.NewLimiter(rate.Limit(cfg.UpLimit), download.LimiterBurst(cfg.UpLimit))
|
|
}
|
|
|
|
// Candidate ports from the --listen-port spec.
|
|
ports := ParsePorts(cfg.ListenPortSpec)
|
|
if len(ports) == 0 {
|
|
// No usable spec: let the OS choose (anacrolix already retries on 0).
|
|
return torrent.NewClient(c)
|
|
}
|
|
|
|
// Try each candidate; the library only auto-retries when the port is 0, so
|
|
// for a fixed busy port it returns an address-in-use error and we advance.
|
|
var lastErr error
|
|
for _, p := range ports {
|
|
c.ListenPort = p
|
|
cl, err := torrent.NewClient(c)
|
|
if err == nil {
|
|
return cl, nil
|
|
}
|
|
if !isAddrInUse(err) {
|
|
return nil, err
|
|
}
|
|
lastErr = err
|
|
}
|
|
// Every named port was busy; fall back to an OS-chosen port rather than
|
|
// giving up on BitTorrent, so the run keeps going with whatever it can bind.
|
|
c.ListenPort = 0
|
|
cl, err := torrent.NewClient(c)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("all listen ports busy (%w); fallback failed: %v", lastErr, err)
|
|
}
|
|
return cl, nil
|
|
}
|
|
|
|
// isAddrInUse reports whether err is an "address already in use" bind failure.
|
|
// (anacrolix/missinggo's IsAddrInUse is this exact string match, so we just do
|
|
// it directly rather than carry the import for it.)
|
|
func isAddrInUse(err error) bool {
|
|
return err != nil && strings.Contains(err.Error(), "address already in use")
|
|
}
|
|
|
|
// Options are the per-download settings.
|
|
type Options struct {
|
|
SeedTimeSet bool // was --seed-time given?
|
|
SeedTime time.Duration // how long to seed (0 with SeedTimeSet means no seeding)
|
|
SeedRatio float64 // stop seeding at this ratio; checked alongside SeedTime
|
|
StopTimeout time.Duration // abort if no download progress for this long (0 = off)
|
|
SelectFiles map[int]bool // 1-based file indexes to fetch; empty = all
|
|
CheckIntegrity bool // re-verify data against piece hashes before downloading
|
|
DryRun bool // fetch metadata only, then stop without downloading (--dry-run)
|
|
}
|
|
|
|
// Download implements download.Download for a single torrent or magnet, added
|
|
// to a shared client.
|
|
type Download struct {
|
|
client *torrent.Client
|
|
source string // magnet URI or path to a .torrent file
|
|
isFile bool
|
|
opts Options
|
|
|
|
name atomic.Pointer[string]
|
|
status int32 // download.Status
|
|
|
|
// t is published once Run adds the torrent to the client. It is atomic so
|
|
// Stat stays lock-free, exactly like httpdl.Stat and bt's own name field.
|
|
t atomic.Pointer[torrent.Torrent]
|
|
}
|
|
|
|
// New builds a torrent download on the shared client. source is a .torrent path
|
|
// when isFile is true, otherwise a magnet URI.
|
|
func New(client *torrent.Client, source string, isFile bool, opts Options) *Download {
|
|
d := &Download{client: client, source: source, isFile: isFile, opts: opts}
|
|
if isFile {
|
|
d.setName(strings.TrimSuffix(filepath.Base(source), ".torrent"))
|
|
} else {
|
|
d.setName(magnetName(source))
|
|
}
|
|
return d
|
|
}
|
|
|
|
func (d *Download) Name() string {
|
|
if p := d.name.Load(); p != nil {
|
|
return *p
|
|
}
|
|
return "torrent"
|
|
}
|
|
|
|
func (d *Download) setName(s string) { d.name.Store(&s) }
|
|
|
|
// setStatus stores the lifecycle state with a single atomic op so Stat stays
|
|
// lock-free, mirroring httpdl. Errored is set once, by Run's deferred guard, when
|
|
// Run returns a non-nil error.
|
|
func (d *Download) setStatus(s download.Status) {
|
|
atomic.StoreInt32(&d.status, int32(s))
|
|
}
|
|
|
|
// id is the process-unique identity for this download: the infohash once known,
|
|
// otherwise the source string (a pre-metadata magnet has no infohash yet, and
|
|
// two such magnets would otherwise collide on Name).
|
|
func (d *Download) id(t *torrent.Torrent) string {
|
|
if t != nil {
|
|
if h := t.InfoHash(); h != (metainfo.Hash{}) {
|
|
return h.HexString()
|
|
}
|
|
}
|
|
return d.source
|
|
}
|
|
|
|
func (d *Download) Stat() download.Stat {
|
|
t := d.t.Load()
|
|
|
|
st := download.Status(atomic.LoadInt32(&d.status))
|
|
if t == nil || t.Info() == nil {
|
|
return download.Stat{Name: d.Name(), ID: d.id(t), IsBT: true, Status: st, Total: -1}
|
|
}
|
|
s := t.Stats()
|
|
total, completed := t.Length(), t.BytesCompleted()
|
|
// With --select-file, report progress against the selected subset, not the
|
|
// whole torrent, so the percentage and ETA mean something.
|
|
if len(d.opts.SelectFiles) > 0 {
|
|
total, completed = 0, 0
|
|
for i, f := range t.Files() {
|
|
if d.opts.SelectFiles[i+1] {
|
|
total += f.Length()
|
|
completed += f.BytesCompleted()
|
|
}
|
|
}
|
|
}
|
|
return download.Stat{
|
|
Name: d.Name(),
|
|
ID: d.id(t),
|
|
IsBT: true,
|
|
Status: st,
|
|
Total: total,
|
|
Completed: completed,
|
|
Uploaded: s.BytesWrittenData.Int64(),
|
|
Conns: s.ActivePeers,
|
|
Seeders: s.ConnectedSeeders,
|
|
}
|
|
}
|
|
|
|
// Run adds the torrent to the shared client, downloads it, then seeds. On any
|
|
// failure the named return is non-nil, and the deferred guard flips the status
|
|
// to Errored exactly once — so each failure site below just returns its error,
|
|
// mirroring httpdl.Run's idiom.
|
|
func (d *Download) Run(ctx context.Context) (err error) {
|
|
defer func() {
|
|
if err != nil {
|
|
d.setStatus(download.Errored)
|
|
}
|
|
}()
|
|
d.setStatus(download.Active)
|
|
if d.client == nil {
|
|
return fmt.Errorf("bittorrent is unavailable")
|
|
}
|
|
|
|
t, err := d.add()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer t.Drop() // release this torrent; the client lives on for others
|
|
d.t.Store(t)
|
|
|
|
if err := d.awaitInfo(ctx, t); err != nil {
|
|
return err
|
|
}
|
|
d.setName(t.Name())
|
|
|
|
// --dry-run: fetching the metadata is all a torrent does under a dry run
|
|
// (the content download is cancelled). Report success without downloading.
|
|
if d.opts.DryRun {
|
|
d.setStatus(download.Complete)
|
|
return nil
|
|
}
|
|
|
|
// --check-integrity: re-hash the existing on-disk data BEFORE arming the
|
|
// request loop, so already-good pieces are not re-requested from peers (aria2
|
|
// verifies first, then fetches only what is missing).
|
|
if d.opts.CheckIntegrity {
|
|
if err := t.VerifyDataContext(ctx); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if err := d.choose(t); err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := d.wait(ctx, t); err != nil {
|
|
return err
|
|
}
|
|
|
|
d.setStatus(download.Seeding)
|
|
d.seed(ctx, t)
|
|
d.setStatus(download.Complete)
|
|
return nil
|
|
}
|
|
|
|
// add hands the source to the shared client. It recovers from the library's
|
|
// panics on malformed input (e.g. a magnet with a zero infohash) so a bad
|
|
// argument fails as an ordinary error instead of crashing the whole program.
|
|
func (d *Download) add() (t *torrent.Torrent, err error) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
t, err = nil, fmt.Errorf("invalid torrent/magnet: %v", r)
|
|
}
|
|
}()
|
|
if d.isFile {
|
|
return d.client.AddTorrentFromFile(d.source)
|
|
}
|
|
return d.client.AddMagnet(d.source)
|
|
}
|
|
|
|
// choose applies --select-file: when set, only the listed files are fetched.
|
|
// With no selection we ask for everything. It is an error for a selection to
|
|
// match no file, so a typo'd index fails loudly instead of "completing" nothing.
|
|
func (d *Download) choose(t *torrent.Torrent) error {
|
|
if len(d.opts.SelectFiles) == 0 {
|
|
t.DownloadAll()
|
|
return nil
|
|
}
|
|
files := t.Files()
|
|
matched := 0
|
|
for i, f := range files {
|
|
if d.opts.SelectFiles[i+1] {
|
|
f.SetPriority(torrent.PiecePriorityNormal)
|
|
matched++
|
|
} else {
|
|
f.SetPriority(torrent.PiecePriorityNone)
|
|
}
|
|
}
|
|
if matched == 0 {
|
|
return fmt.Errorf("--select-file: no file index in 1..%d", len(files))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// metadataTimeout bounds the magnet metadata-fetch phase. A magnet with no
|
|
// reachable peers — a dead link, or UDP trackers behind a firewall with no DHT —
|
|
// would otherwise hang forever, so awaitInfo gives up after this fixed window.
|
|
// It is not a user-facing flag: aria2 has no --bt-metadata-timeout, and 60s is a
|
|
// generous ceiling that a healthy swarm clears in well under a second.
|
|
const metadataTimeout = 60 * time.Second
|
|
|
|
// awaitInfo blocks until the torrent metadata arrives, giving up after
|
|
// metadataTimeout so a peerless magnet fails fast instead of hanging forever.
|
|
func (d *Download) awaitInfo(ctx context.Context, t *torrent.Torrent) error {
|
|
tm := time.NewTimer(metadataTimeout)
|
|
defer tm.Stop()
|
|
select {
|
|
case <-t.GotInfo():
|
|
return nil
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-tm.C:
|
|
// Wrap DeadlineExceeded so the exit-code mapping classifies it as a timeout
|
|
// (2) via errors.Is, rather than matching on the message text.
|
|
return fmt.Errorf("timed out fetching metadata after %s: %w", metadataTimeout, context.DeadlineExceeded)
|
|
}
|
|
}
|
|
|
|
// wait blocks until the wanted data is complete, ctx is cancelled, or the
|
|
// download stalls past --bt-stop-timeout.
|
|
func (d *Download) wait(ctx context.Context, t *torrent.Torrent) error {
|
|
tick := time.NewTicker(500 * time.Millisecond)
|
|
defer tick.Stop()
|
|
|
|
last := t.BytesCompleted()
|
|
lastChange := time.Now()
|
|
for {
|
|
if d.complete(t) {
|
|
return nil
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-tick.C:
|
|
if n := t.BytesCompleted(); n != last {
|
|
last, lastChange = n, time.Now()
|
|
} else if d.opts.StopTimeout > 0 && time.Since(lastChange) > d.opts.StopTimeout {
|
|
return fmt.Errorf("no progress for %s", d.opts.StopTimeout)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// complete reports whether everything we asked for is downloaded AND verified.
|
|
// For a selection we check each wanted file's pieces are hash-complete rather
|
|
// than trusting byte counts, which include not-yet-verified (dirty) bytes.
|
|
func (d *Download) complete(t *torrent.Torrent) bool {
|
|
if len(d.opts.SelectFiles) == 0 {
|
|
return t.Complete().Bool()
|
|
}
|
|
for i, f := range t.Files() {
|
|
if !d.opts.SelectFiles[i+1] {
|
|
continue
|
|
}
|
|
for _, ps := range f.State() {
|
|
if !ps.Complete {
|
|
return false
|
|
}
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
// seed keeps the torrent up after completion. Seeding stops as soon as EITHER
|
|
// criterion fires: the --seed-time elapses or the --seed-ratio is reached. Both
|
|
// are active at once when both are set; with neither set we seed forever (until
|
|
// ctx is cancelled). One short-circuit applies: --seed-time=0 means do not seed.
|
|
func (d *Download) seed(ctx context.Context, t *torrent.Torrent) {
|
|
if d.opts.SeedTimeSet && d.opts.SeedTime == 0 {
|
|
return
|
|
}
|
|
// Arm the time criterion whenever --seed-time was given (independently of the
|
|
// ratio criterion below).
|
|
var timer <-chan time.Time
|
|
if d.opts.SeedTimeSet {
|
|
tm := time.NewTimer(d.opts.SeedTime)
|
|
defer tm.Stop()
|
|
timer = tm.C
|
|
}
|
|
tick := time.NewTicker(time.Second)
|
|
defer tick.Stop()
|
|
|
|
// Share ratio is uploaded/downloaded. The denominator is the torrent's
|
|
// completed bytes at the moment seeding starts — under --select-file only the
|
|
// selected files are ever fetched, so this is already close to the subset, but
|
|
// it is not subset-filtered: any boundary/shared pieces count toward it too.
|
|
dl := t.BytesCompleted()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-timer:
|
|
return
|
|
case <-tick.C:
|
|
// Ratio criterion is checked every tick whenever a ratio is set,
|
|
// regardless of whether --seed-time is also set.
|
|
if d.opts.SeedRatio > 0 && dl > 0 {
|
|
st := t.Stats()
|
|
up := st.BytesWrittenData.Int64()
|
|
if float64(up)/float64(dl) >= d.opts.SeedRatio {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// FileInfo describes one file inside a torrent, for --show-files.
|
|
type FileInfo struct {
|
|
Index int
|
|
Path string
|
|
Length int64
|
|
}
|
|
|
|
// Files lists the files in a .torrent without starting a download.
|
|
func Files(torrentPath string) ([]FileInfo, error) {
|
|
mi, err := metainfo.LoadFromFile(torrentPath)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
info, err := mi.UnmarshalInfo()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if len(info.Files) == 0 {
|
|
return []FileInfo{{1, info.Name, info.Length}}, nil
|
|
}
|
|
out := make([]FileInfo, len(info.Files))
|
|
for i, f := range info.Files {
|
|
out[i] = FileInfo{i + 1, strings.Join(f.Path, "/"), f.Length}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// SourceInfoHash returns the v1 infohash of a .torrent file (isFile) or a magnet
|
|
// URI, with ok=false when it can't be known without a network fetch (an HTTP
|
|
// URL), the source is malformed, or the v1 infohash is zero — so a BitTorrent
|
|
// v2-only source is treated as unique (not deduped). Callers use it to collapse
|
|
// two sources that name the same torrent into one download: the shared client
|
|
// keys torrents by infohash and does not refcount, so two jobs sharing one
|
|
// torrent would have the first to finish Drop it out from under the other.
|
|
func SourceInfoHash(source string, isFile bool) (metainfo.Hash, bool) {
|
|
var h metainfo.Hash
|
|
if isFile {
|
|
mi, err := metainfo.LoadFromFile(source)
|
|
if err != nil {
|
|
return metainfo.Hash{}, false
|
|
}
|
|
h = mi.HashInfoBytes()
|
|
} else {
|
|
m, err := metainfo.ParseMagnetUri(source)
|
|
if err != nil {
|
|
return metainfo.Hash{}, false
|
|
}
|
|
h = m.InfoHash
|
|
}
|
|
if h == (metainfo.Hash{}) {
|
|
return metainfo.Hash{}, false
|
|
}
|
|
return h, true
|
|
}
|
|
|
|
func magnetName(uri string) string {
|
|
if m, err := metainfo.ParseMagnetUri(uri); err == nil && m.DisplayName != "" {
|
|
return m.DisplayName
|
|
}
|
|
return "magnet"
|
|
}
|