Files
got/bt/bt.go
Hojun-Cho 10672f3816 all: cut bloat found in the pike audit
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.
2026-06-21 15:45:56 +09:00

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"
}