Files
jellybit/internal/worker/worker.go
T
2026-07-17 22:11:53 +03:00

1152 lines
63 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package worker — владелец машины состояний. Поллит qBittorrent по
// категории, переводит задачи между состояниями и сериализует команды
// транспортов (cancel/retry), чтобы два транспорта не гонялись за одно
// состояние.
//
// Ф1 ведёт задачу downloading → completed, плюс stuck/failed по таймаутам и
// ошибкам qBittorrent. Ф3 продолжает: completed → recognizing (вызов
// recognize) → review; команды ревью (apply/refine/reject/defer/undo,
// переключение типа, пометка «игнор») раскладывают файлы хардлинками через
// layout. Распознавание зовётся в поллинг-цикле, команды — из транспортов;
// всё под per-download блокировкой w.mu.
package worker
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"slices"
"strings"
"sync"
"time"
"git.vakhrushev.me/av/jellybit/internal/ident"
"git.vakhrushev.me/av/jellybit/internal/layout"
"git.vakhrushev.me/av/jellybit/internal/logctx"
"git.vakhrushev.me/av/jellybit/internal/magnet"
"git.vakhrushev.me/av/jellybit/internal/naming"
"git.vakhrushev.me/av/jellybit/internal/qbt"
"git.vakhrushev.me/av/jellybit/internal/recognize"
"git.vakhrushev.me/av/jellybit/internal/store"
"git.vakhrushev.me/av/jellybit/internal/torrent"
)
// Стадии (capability) — адресуют запись к подсистеме при корреляции по
// download_id. Уровень логов от стадии не зависит.
const (
capIngest = "ingest" // приём/скачивание (поллинг, reconcile, переходы)
capRecognize = "recognition" // распознавание фильма/сериала
capFileLayout = "file-layout" // раскладка хардлинками
capReview = "review" // ручные команды ревью
)
// Store — нужная worker часть хранилища.
type Store interface {
ListDownloadsByState(ctx context.Context, states ...store.State) ([]store.Download, error)
ListRecoverable(ctx context.Context, codes ...string) ([]store.Download, error)
GetDownload(ctx context.Context, id string) (*store.Download, error)
SetDownloadState(ctx context.Context, id string, state store.State, errCode, errMsg string) error
// PromoteCatched атомарно переводит catched → downloading с записью имени
// (гард state='catched' — ре-валидация после сетевых вызовов вне блокировки).
PromoteCatched(ctx context.Context, id, displayName string) error
// SetDisplayName обновляет отображаемое имя постфактум (перелив канонического
// имени после распознавания) — без гарда состояния, FSM не двигает.
SetDisplayName(ctx context.Context, id, name string) error
SetParsedContext(ctx context.Context, id, jsonStr string) error
SetSourceMissCount(ctx context.Context, id string, n int) error
SetSourceAddedAt(ctx context.Context, id string, t time.Time) error
// SetRetriedAt проставляет время ручного retry — сброс базиса отсчёта
// таймаутов (magnet_timeout/stuck_after), чтобы возврат в downloading не
// ронял задачу снова на ближайшем тике.
SetRetriedAt(ctx context.Context, id string, t time.Time) error
// Идентичность/инвариант «одна активная загрузка на infohash».
ExistsByInfohash(ctx context.Context, hashes ...string) (bool, error)
CreateDownloadIfNoActive(ctx context.Context, d *store.Download, hashes []string, torrentBlob []byte) (*store.Download, error)
// GetTorrentData — сохранённые байты `.torrent` (для добавления раздачи
// файлом на шаге processCatched/Retry у source_type=torrent).
GetTorrentData(ctx context.Context, downloadID string) ([]byte, error)
ActivateIfNoOtherActive(ctx context.Context, id string, state store.State, errCode, errMsg string) error
AddInfohashes(ctx context.Context, downloadID string, hashes []string) error
// Ф3: распознавание, ревью, раскладка.
CreateRecognition(ctx context.Context, r *store.Recognition, reasons []string) (string, error)
GetCurrentRecognition(ctx context.Context, downloadID string) (*store.Recognition, error)
AddHint(ctx context.Context, downloadID string, text string) error
ListHints(ctx context.Context, downloadID string) ([]string, error)
SetOverride(ctx context.Context, downloadID string, field, value string) error
ListOverrides(ctx context.Context, downloadID string) (map[string]string, error)
CreateFileLinks(ctx context.Context, links []store.FileLink) error
SupersedeForeignLinks(ctx context.Context, downloadID string, dstPaths []string) error
// LiveTitleFolders — dst_path живых ссылок загрузок того же (provider,
// provider_id), кроме excludeDownloadID (правило сходимости папки).
LiveTitleFolders(ctx context.Context, provider, providerID, excludeDownloadID string) ([]string, error)
LatestBatchID(ctx context.Context, downloadID string) (string, error)
ListFileLinksByBatch(ctx context.Context, batchID string) ([]store.FileLink, error)
DeleteFileLinksByBatch(ctx context.Context, batchID string) error
// Кандидаты базы метаданных (ручной выбор в review).
CreateCandidates(ctx context.Context, cands []store.MetadataCandidate) error
ListCandidatesByRecognition(ctx context.Context, recognitionID string) ([]store.MetadataCandidate, error)
GetCandidate(ctx context.Context, id string) (*store.MetadataCandidate, error)
SetCandidateChosen(ctx context.Context, recognitionID, candidateID string) error
}
// QBittorrent — нужная worker часть клиента qBittorrent.
type QBittorrent interface {
Torrents(ctx context.Context, category string) ([]qbt.Torrent, error)
Add(ctx context.Context, ar qbt.AddRequest) error
Files(ctx context.Context, hash string) ([]qbt.File, error)
Delete(ctx context.Context, hashes []string, deleteFiles bool) error
// RenameTorrent задаёт имя уже добавленной раздачи по её ключу (Torrent.Hash).
RenameTorrent(ctx context.Context, hash, name string) error
}
// Recognizer — распознаватель (recognize.Recognizer).
type Recognizer interface {
Recognize(ctx context.Context, in recognize.Input) (recognize.Result, error)
// Director тянет режиссёра выбранного источника по (provider, id) из
// метабазы (credits) — для ручного выбора кандидата в ревью. Best-effort:
// нет провайдера/не умеет/ошибка → пустая строка.
Director(ctx context.Context, mt recognize.MediaType, provider, providerID string) string
}
// Namer выводит человекочитаемое отображаемое имя из контекста (naming.Namer).
// Derive возвращает имя (пустое → rename в qBittorrent не задаём) и извлечённую
// структуру (её JSON сохраняется как parsed_context — базовый слой полей имени).
// nil-namer → имя не выводим.
type Namer interface {
Derive(ctx context.Context, contextText, hint string) (name string, fields naming.Fields, ok bool)
}
// Layouter — раскладчик хардлинками (layout.Layouter).
type Layouter interface {
BuildLinks(p layout.Plan) ([]layout.Link, error)
Apply(ctx context.Context, links []layout.Link) ([]layout.Result, error)
Undo(ctx context.Context, links []layout.Link) (int, error)
Remove(ctx context.Context, links []layout.Link) (int, error)
// TitleFolder разбирает dst_path в папку тайтла и базу имени (правило
// сходимости папки). ok=false, если путь не под корнем нужной библиотеки.
TitleFolder(t layout.MediaType, dst string) (dir, base string, ok bool)
}
// NotifyEvent — повод позвать пользователя.
type NotifyEvent string
const (
EventReview NotifyEvent = "review" // задача ждёт подтверждения
EventDone NotifyEvent = "done" // раскладка завершена
EventOrphaned NotifyEvent = "orphaned" // источник пропал, цель — последняя копия
EventTargetMissing NotifyEvent = "target_missing" // цель удалена, доступен relink
EventFailed NotifyEvent = "failed" // задача упала/зависла (failed/stuck)
)
// Коды ошибок (error_code) при переходе в failed/stuck. Восстановимые
// (magnet_timeout/stalled) — следствие нашей нетерпеливости: сверка воскрешает
// такие задачи при оживлении источника (см. reconcileRecovery). qbit_error —
// реальная ошибка qBittorrent, восстановлению не подлежит.
const (
errCodeMagnetTimeout = "magnet_timeout"
errCodeStalled = "stalled"
errCodeQbitError = "qbit_error"
// errCodeQbitAdd — не удалось добавить пойманную загрузку в qBittorrent за
// catch_timeout (устойчивая недоступность qBit). Раздачи в qBittorrent нет,
// восстановлению сверкой не подлежит.
errCodeQbitAdd = "qbit_add"
// errCodeSourceGone — раздача активной (downloading) загрузки устойчиво (после
// дебаунса source_missing_threshold) пропала из qBittorrent: пользователь/другой
// клиент её удалил. Distinct-код, отличный от qbit_error (реальная ошибка qBit) и
// magnet_timeout/stalled (наша нетерпеливость). Восстановлению сверкой НЕ подлежит
// (удаление намеренно) — но задача штатно retriable: Retry заново отдаёт источник.
errCodeSourceGone = "source_gone"
)
// Notifier — исходящие пинги (Telegram). Вызывается неблокирующе.
type Notifier interface {
Notify(ctx context.Context, downloadID string, event NotifyEvent)
}
// Scanner — триггер пересканирования медиатеки Jellyfin. Вызывается
// неблокирующе после успешной раскладки, чтобы новые файлы быстрее появились
// в проигрывателе.
type Scanner interface {
RefreshLibraries(ctx context.Context) error
}
// Config — параметры воркера.
type Config struct {
Category string
Tag string // метка для усыновления существующих раздач (discovery)
SavePath string
PathMap map[string]string // трансляция save_path qBit → хост-путь (обычно пусто)
PollInterval time.Duration
StuckAfter time.Duration // stalledDL дольше → stuck
MagnetTimeout time.Duration // metaDL дольше → failed
CatchTimeout time.Duration // catched дольше (не удалось добавить в qBit) → failed
// SourceMissingThreshold — порог дебаунса пропажи источника (тиков сверки).
// <1 трактуется как 1 (помечаем при первой же устойчивой пропаже).
SourceMissingThreshold int
}
// Live — живая телеметрия одной раздачи из снимка воркера. Курированный срез
// qbt.Torrent: читатели (httpapi) не зависят от пакета qbt, контракт чтения
// узкий. Seeding вычисляется воркером (classify), чтобы трактовка завершённости
// не дублировалась в транспорте.
type Live struct {
Progress float64 // доля 0..1
DlSpeed int64 // скорость загрузки, байт/с
ETA int64 // оценка до завершения, с (8640000 ≈ ∞)
State string // сырое состояние qBittorrent
Seeding bool // торрент завершён и раздаётся
TotalSize int64 // полный размер раздачи, байт (доступен для любой раздачи в снимке)
Ratio float64 // рейтинг отдачи (может быть <0 = ∞/н/д)
Seeds int // подключённые сиды
Peers int // подключённые личи
Uploaded int64 // отдано всего, байт
UpSpeed int64 // скорость отдачи, байт/с
}
// liveFrom собирает Live из торрента qBittorrent (Seeding — через classify).
func liveFrom(t qbt.Torrent) Live {
return Live{
Progress: t.Progress,
DlSpeed: t.Dlspeed,
ETA: t.Eta,
State: t.State,
Seeding: classify(t.State) == classReady,
TotalSize: t.TotalSize,
Ratio: t.Ratio,
Seeds: t.NumSeeds,
Peers: t.NumLeechs,
Uploaded: t.Uploaded,
UpSpeed: t.Upspeed,
}
}
// Worker — поллер и владелец переходов.
type Worker struct {
store Store
qbt QBittorrent
recognizer Recognizer
layouter Layouter
namer Namer // опц. вывод отображаемого имени на шаге добавления catched
cfg Config
log *slog.Logger
mu sync.Mutex // сериализует переходы (поллинг + команды)
now func() time.Time // подменяется в тестах
newID func() string // генератор apply_batch_id (подменяется в тестах)
notifier Notifier // опц. исходящие пинги
scanner Scanner // опц. пересканирование Jellyfin
// live — снимок живой телеметрии раздач (ключ — lowercase infohash, по три
// ключа на торрент, как byHash). Обновляется атомарным свопом карты на
// каждом тике Poll. Отдельный RWMutex (не w.mu): UI читает телеметрию часто,
// смешивать частые чтения с замком переходов — лишняя конкуренция. Снимок
// волатилен, в БД не хранится.
liveMu sync.RWMutex
live map[string]Live
// failNotified — дебаунс повторных EventFailed по задаче (download_id →
// время последнего пинга). Мерцающий stalled-торрент колеблется
// stuck↔downloading; без дебаунса каждый цикл слал бы уведомление. Память
// процесса: при рестарте дебаунс сбрасывается — допустимо. Доступ под w.mu.
failNotified map[string]time.Time
}
// failNotifyDebounce — минимальный интервал между уведомлениями о падении
// одной задачи (см. failNotified).
const failNotifyDebounce = time.Hour
// SetNamer подключает вывод отображаемого имени для шага добавления catched
// (до запуска Run). nil → имя не выводим, добавляем без rename.
func (w *Worker) SetNamer(n Namer) { w.namer = n }
// SetNotifier подключает исходящие пинги (до запуска Run).
func (w *Worker) SetNotifier(n Notifier) { w.notifier = n }
// SetScanner подключает пересканирование Jellyfin (до запуска Run).
func (w *Worker) SetScanner(s Scanner) { w.scanner = s }
// New собирает воркер. recognizer/layouter могут быть nil (Ф1 без Ф3-ступеней
// распознавания и раскладки) — тогда completed-задачи не двигаются дальше.
func New(st Store, qb QBittorrent, rec Recognizer, lay Layouter, cfg Config, log *slog.Logger) *Worker {
return &Worker{
store: st,
qbt: qb,
recognizer: rec,
layouter: lay,
cfg: cfg,
log: log,
now: time.Now,
newID: defaultBatchID,
failNotified: map[string]time.Time{},
live: map[string]Live{},
}
}
// Live возвращает живую телеметрию раздачи по infohash (любому из v1/v2/hash).
// ok=false, если infohash пуст или раздачи не было в последнем тике поллинга —
// тогда читатель деградирует без живых значений. Чтение под RLock.
func (w *Worker) Live(infohash string) (Live, bool) {
if infohash == "" {
return Live{}, false
}
w.liveMu.RLock()
defer w.liveMu.RUnlock()
l, ok := w.live[strings.ToLower(infohash)]
return l, ok
}
// setLive атомарно подменяет снимок телеметрии готовой картой.
func (w *Worker) setLive(snap map[string]Live) {
w.liveMu.Lock()
w.live = snap
w.liveMu.Unlock()
}
// defaultBatchID — идентификатор батча раскладки (ULID, единая точка
// генерации id — internal/ident; сортируем по времени, удобен в логах).
func defaultBatchID() string {
return ident.NewID()
}
// scoped кладёт в ctx scoped-логгер загрузки (capability + download_id
// [+ infohash]); стадии и внешние клиенты достают его из ctx и дописывают эти
// ключи на каждую запись сами — без ручного доклеивания download_id.
func (w *Worker) scoped(ctx context.Context, capability string, id string, infohash string) context.Context {
log := w.log.With("capability", capability, "download_id", id)
if infohash != "" {
log = log.With("infohash", infohash)
}
return logctx.With(ctx, log)
}
// Run крутит цикл поллинга до отмены ctx.
func (w *Worker) Run(ctx context.Context) {
w.log.Info("worker started", "poll_interval", w.cfg.PollInterval, "category", w.cfg.Category)
t := time.NewTicker(w.cfg.PollInterval)
defer t.Stop()
w.pollOnce(ctx)
for {
select {
case <-ctx.Done():
w.log.Info("worker stopped")
return
case <-t.C:
w.pollOnce(ctx)
}
}
}
func (w *Worker) pollOnce(ctx context.Context) {
if err := w.Poll(ctx); err != nil {
w.log.Warn("poll failed", "error", err)
}
// Быстрый приём отложил добавление в qBittorrent: подхватываем пойманные
// (catched) загрузки и добавляем их (сеть — вне блокировки переходов).
w.processCatched(ctx)
// Восстанавливаем задачи, застрявшие в linking после краха между claim и
// финальным переходом (иначе их не листит никто — вечный лимбо).
w.sweepLinking(ctx)
// Ф3: распознаём завершённые загрузки (и перезапускаем по подсказке).
if w.recognizer != nil {
w.recognizePending(ctx)
}
}
// sweepLinking восстанавливает задачи, застрявшие в состоянии linking. Любая
// linking-задача, видимая под w.mu, устарела по построению: активная раскладка
// (linkPlan) держит w.mu на всё время и завершает переход из linking ДО отпускания
// замка — значит эта задача осталась в linking после краха процесса между claim
// (переходом в linking) и финальным переходом. Возвращаем её в review с причиной;
// человек повторит Apply (linkPlan идемпотентен), и незаписанный учёт хардлинков
// допишется. Так у linking появляется владелец на рестарте/тике — инвариант «у
// каждого нетерминального состояния есть владелец» (как recognizePending для
// recognizing). Выполняется на каждом тике и на старте (первый pollOnce до цикла).
func (w *Worker) sweepLinking(ctx context.Context) {
w.mu.Lock()
defer w.mu.Unlock()
stuck, err := w.store.ListDownloadsByState(ctx, store.StateLinking)
if err != nil {
w.log.Warn("sweep linking list failed", "capability", capFileLayout, "error", err)
return
}
for _, d := range stuck {
lctx := w.scoped(ctx, capFileLayout, d.ID, d.PrimaryInfohash())
w.transition(lctx, d, store.StateReview, "interrupted",
"прерванная раскладка, повтори применение")
}
}
// processCatched — асинхронный шаг добавления пойманных загрузок в qBittorrent.
// Для каждой catched: (предохранитель) если висит дольше catch_timeout — уводим
// в failed; иначе выводим имя и добавляем в qBit. Медленные вызовы (LLM-namer,
// qbt.Add) идут ВНЕ w.mu, чтобы не задерживать команды транспортов и поллинг;
// под w.mu берутся только короткие DB-переходы (с ре-валидацией state=catched).
func (w *Worker) processCatched(ctx context.Context) {
w.mu.Lock()
catched, err := w.store.ListDownloadsByState(ctx, store.StateCatched)
w.mu.Unlock()
if err != nil {
w.log.Warn("list catched failed", "capability", capIngest, "error", err)
return
}
if len(catched) == 0 {
return
}
// Снимок присутствия раздач в qBittorrent (один листинг на тик): по нему ДО
// вызова namer решаем, добавлять ли задачу вообще. Провал листинга —
// qBittorrent недоступен: пойманные в этот тик не трогаем (ни namer, ни Add),
// повтор на следующем; устойчивая недоступность отсекается предохранителем
// catch_timeout.
torrents, err := w.qbt.Torrents(ctx, "")
if err != nil {
w.log.Warn("list torrents for catched failed", "capability", capIngest, "error", err)
return
}
byHash := torrentsByHash(torrents)
for _, d := range catched {
cctx := w.scoped(ctx, capIngest, d.ID, d.PrimaryInfohash())
// Торрент уже в qBittorrent — повторный Add не нужен (и вреден: qBittorrent
// отверг бы дубль, 409) и LLM-namer не зовём: усыновляем раздачу, доводя
// задачу до downloading. См. promoteExisting.
if t, ok := torrentFor(d, byHash); ok {
w.promoteExisting(cctx, d, t)
continue
}
// Предохранитель: устойчивая невозможность добавить в qBittorrent (раздачи
// в снимке нет и висит дольше catch_timeout).
if w.cfg.CatchTimeout > 0 {
if age, ok := w.catchedAge(d); ok && age > w.cfg.CatchTimeout {
w.mu.Lock()
// Ре-валидация под замком: список catched снят раньше, задачу
// могли отменить (catched → cancelled) в это окно — тогда failed
// не навязываем (иначе затёрли бы cancelled и слали лишний пинг).
if cur, err := w.store.GetDownload(cctx, d.ID); err == nil && cur.State == store.StateCatched {
w.transition(cctx, d, store.StateFailed, errCodeQbitAdd,
fmt.Sprintf("not added to qBittorrent after %s", age.Truncate(time.Second)))
}
w.mu.Unlock()
continue
}
}
// Раздачи в qBittorrent нет — обычный путь добавления. Перечитываем запись
// под замком: (а) актуальный source_type (апгрейд magnet→torrent мог
// случиться после снятия списка catched — иначе добавили бы magnet из
// устаревшего снимка), (б) ре-валидация state=catched. Тяжёлые вызовы
// (чтение байтов, namer, Add) — вне замка.
w.mu.Lock()
cur, gerr := w.store.GetDownload(cctx, d.ID)
fresh := gerr == nil && cur != nil && cur.State == store.StateCatched
w.mu.Unlock()
if !fresh {
continue // отменили/пропала, пока шёл листинг — не трогаем
}
hint, addReq, prepErr := w.sourceAddParts(cctx, *cur)
if prepErr != nil {
// Байты torrent недоступны (не должно быть при штатном приёме) —
// остаёмся в catched, повтор на следующем тике.
logctx.From(cctx).Warn("catched prepare add failed, will retry", "error", prepErr)
continue
}
var rename string
if w.namer != nil {
name, fields, extracted := w.namer.Derive(cctx, cur.Context, hint)
rename = name
// Сохраняем извлечённую структуру как базовый слой полей имени
// (parsed_context) — режиссёр/год из контекста не теряются и
// переиспользуются при обновлении display_name. Best-effort: сбой не
// валит добавление (косметика). Пустая структура → нечего сохранять.
if extracted {
if blob, mErr := json.Marshal(fields); mErr == nil {
if sErr := w.store.SetParsedContext(cctx, cur.ID, string(blob)); sErr != nil {
logctx.From(cctx).Warn("catched save parsed context failed", "error", sErr)
}
}
}
}
addReq.Rename = rename
// F3: re-read состояния прямо перед Add — вывод имени (LLM) шёл секунды вне
// замка, задачу могли отменить (catched → cancelled). Если уже не catched,
// источник в qBittorrent не добавляем вовсе (иначе остался бы неуправляемый
// торрент без задачи-владельца).
w.mu.Lock()
before, berr := w.store.GetDownload(cctx, d.ID)
stillCatched := berr == nil && before != nil && before.State == store.StateCatched
w.mu.Unlock()
if !stillCatched {
logctx.From(cctx).Info("catched add skipped before qbittorrent", "reason", "no longer catched")
continue
}
// F3, гарантия «удаляем только своё»: свежим листингом (вне замка)
// подтверждаем, что раздачи с нашим infohash в qBittorrent ЕЩЁ НЕТ. Только
// тогда торрент, появившийся под этим хешем сразу после нашего Add, — наш
// артефакт, и позднейшая уборка вправе снести его С ДАННЫМИ. Сбой листинга →
// отсутствие не подтверждено, Add не делаем (повтор на следующем тике),
// иначе delete-с-данными стал бы небезопасен. Присутствие → внешний клиент
// добавил тот же торрент в окно гонки: Add не делаем, усыновит следующий тик
// (promoteExisting); чужие данные не трогаем.
snap, ferr := w.qbt.Torrents(cctx, "")
if ferr != nil {
logctx.From(cctx).Warn("catched presence recheck failed, will retry", "error", ferr)
continue
}
if _, present := torrentFor(*before, torrentsByHash(snap)); present {
logctx.From(cctx).Info("catched torrent already present in qbittorrent, will adopt")
continue
}
addErr := w.qbt.Add(cctx, addReq)
if addErr != nil {
// Транзиентный сбой (qBit отверг/недоступен) — остаёмся в catched,
// повтор на следующем тике. Вызов qBit уже залогировал клиент (ext.*).
logctx.From(cctx).Warn("catched add to qbittorrent failed, will retry", "error", addErr)
continue
}
// Успех Add: короткий переход под w.mu с ре-валидацией state=catched
// (загрузку могли отменить, пока шли сетевые вызовы).
w.mu.Lock()
perr := w.store.PromoteCatched(cctx, d.ID, rename)
if perr == nil {
logctx.From(cctx).Info("state transition", "from", store.StateCatched,
"to", store.StateDownloading)
w.mu.Unlock()
continue
}
// Промоут не прошёл. Причину определяем СВЕЖИМ состоянием под тем же замком,
// а НЕ текстом ошибки (см. errors.md): PromoteCatched возвращает ошибку и при
// отмене (state != catched), и при транзиентном сбое БД (state всё ещё
// catched, задача жива).
after, aerr := w.store.GetDownload(cctx, d.ID)
cancelled := aerr == nil && after != nil && after.State != store.StateCatched
w.mu.Unlock()
if !cancelled {
// Транзиентный сбой БД (или не смогли перечитать) — торрент наш и живой,
// не удаляем: переход доведётся на следующем тике усыновлением
// присутствующей раздачи (promoteExisting).
logctx.From(cctx).Warn("catched promote failed, will retry", "error", perr)
continue
}
// F3: отмена (catched → cancelled) в окне между Add и записью перехода.
// Торрент добавлен НАМИ этим Add (отсутствие infohash подтверждено выше), а
// задачи-владельца больше нет — снимаем свой артефакт С ДАННЫМИ. Инвариант
// «источник неприкосновенен» защищает пользовательские данные, а не наш
// только что добавленный торрент; удаление идёт через API qBittorrent, не
// прямыми fs-операциями.
logctx.From(cctx).Warn("torrent left in qbittorrent after cancel, removing", "error", perr)
if delErr := w.qbt.Delete(cctx, before.HashList(), true); delErr != nil {
logctx.From(cctx).Error("cleanup added torrent after cancel failed", "error", delErr)
} else {
logctx.From(cctx).Warn("added torrent removed after cancel")
}
}
}
// promoteExisting усыновляет пойманную загрузку, чей торрент уже присутствует в
// qBittorrent (снимок тика): переводит catched → downloading БЕЗ повторного Add
// (иначе qBittorrent отверг бы дубль — 409 — и задача зациклилась бы) и без LLM.
// Имя берём из раздачи снимка (t.Name); у свежего magnet без метаданных (metaDL)
// оно может быть пустым — распознавание дольёт имя позже. Инвариант приёма
// гарантирует, что сюда доходит лишь загрузка без другой активной задачи на тот
// же infohash, поэтому присутствие раздачи трактуем как «усыновить и разложить»,
// а не как конфликт. Короткий DB-переход под w.mu; атомарный гард PromoteCatched
// (WHERE state='catched') сам отсекает гонку отмены, случившуюся, пока шёл листинг
// вне замка, — отдельный re-read не нужен.
func (w *Worker) promoteExisting(ctx context.Context, d store.Download, t qbt.Torrent) {
w.mu.Lock()
defer w.mu.Unlock()
if err := w.store.PromoteCatched(ctx, d.ID, t.Name); err != nil {
logctx.From(ctx).Info("catched promote skipped", "reason", err.Error())
return
}
logctx.From(ctx).Info("state transition", "from", store.StateCatched,
"to", store.StateDownloading, "reason", "already present in qbittorrent")
}
// torrentsByHash индексирует раздачи по каждому из их хешей (lowercase), как это
// делает Poll для своего снимка. Использует processCatched, чтобы проверить, есть
// ли торрент пойманной загрузки уже в qBittorrent.
func torrentsByHash(torrents []qbt.Torrent) map[string]qbt.Torrent {
byHash := make(map[string]qbt.Torrent, len(torrents)*2)
for _, t := range torrents {
for _, h := range torrentIndexHashes(t) {
byHash[strings.ToLower(h)] = t
}
}
return byHash
}
// catchedAge — возраст пойманной загрузки от created_at (у catched раздачи в
// qBittorrent ещё нет, added_on недоступен). ok=false — created_at не разобрать.
func (w *Worker) catchedAge(d store.Download) (time.Duration, bool) {
created, err := d.CreatedTime()
if err != nil {
w.log.Warn("cannot determine catched age",
"capability", capIngest, "download_id", d.ID,
"created_at", d.CreatedAt, "error", err)
return 0, false
}
return w.now().Sub(created), true
}
// sourceAddParts собирает параметры добавления раздачи в qBittorrent по типу
// источника и подсказку имени (hint) для namer. magnet/url — ссылкой (URLs),
// hint из полей ссылки; torrent — байтами файла (Torrents), hint из имени
// раздачи. Общий для обоих add-путей воркера (processCatched и Retry), чтобы
// диспетч по source_type был в одном месте. Rename вызывающий проставляет сам
// (после namer). Ошибку возвращает лишь torrent-ветка (байты недоступны).
func (w *Worker) sourceAddParts(ctx context.Context, d store.Download) (hint string, req qbt.AddRequest, err error) {
req = qbt.AddRequest{Category: w.cfg.Category, SavePath: w.cfg.SavePath}
if d.SourceType == store.SourceTorrent {
data, derr := w.store.GetTorrentData(ctx, d.ID)
if derr != nil {
return "", qbt.AddRequest{}, fmt.Errorf("torrent bytes: %w", derr)
}
if info, perr := torrent.Parse(data); perr == nil {
hint = info.DisplayName
}
req.Torrents = [][]byte{data}
return hint, req, nil
}
// magnet/url: SourceRef — добавляемая ссылка, и источник hint для namer.
if info, perr := magnet.Parse(d.SourceRef); perr == nil {
hint = info.DisplayName
}
req.URLs = []string{d.SourceRef}
return hint, req, nil
}
// Poll сверяет активные задачи с состоянием qBittorrent и двигает их.
// Листаем все торренты (а не только свою категорию), чтобы reconcile нашёл и
// усыновлённые по тегу раздачи, а discovery — увидел новые.
func (w *Worker) Poll(ctx context.Context) error {
torrents, err := w.qbt.Torrents(ctx, "")
if err != nil {
return fmt.Errorf("poll: list torrents: %w", err)
}
byHash := make(map[string]qbt.Torrent, len(torrents)*2)
live := make(map[string]Live, len(torrents)*2)
for _, t := range torrents {
l := liveFrom(t)
for _, h := range torrentIndexHashes(t) {
key := strings.ToLower(h)
byHash[key] = t
live[key] = l
}
}
// Снимок зависит только от torrents — свопаем сразу, до store-операций
// ниже (их ранний return по ошибке не должен лишать UI свежей телеметрии).
w.setLive(live)
w.mu.Lock()
defer w.mu.Unlock()
// Усыновляем новые раздачи с нашей категорией/тегом до reconcile.
w.discover(ctx, torrents)
active, err := w.store.ListDownloadsByState(ctx, store.StateDownloading)
if err != nil {
return fmt.Errorf("poll: list active: %w", err)
}
for _, d := range active {
if len(d.Infohashes) == 0 {
continue // нечем сопоставить (в Ф1 не случается: magnet всегда с infohash)
}
t, ok := torrentFor(d, byHash)
// Дебаунс пропажи источника у активной загрузки (тот же счётчик, что и
// сверка рассинхрона): промах наращивает source_miss_count, появление
// раздачи сбрасывает его. Устойчивая пропажа (после порога) уводит задачу
// в failed(source_gone) — иначе downloading без раздачи в qBittorrent
// оставался бы вечным зомби (MAJOR-3).
if present := w.debounceSource(ctx, d, ok); !present {
lctx := w.scoped(ctx, capIngest, d.ID, d.PrimaryInfohash())
logctx.From(lctx).Warn("active download source gone from qbittorrent",
"miss_count", d.SourceMissCount+1)
w.transition(lctx, d, store.StateFailed, errCodeSourceGone,
"источник удалён из qBittorrent")
continue
}
if !ok {
// До порога: транзиентный промах (например рестарт демона qBit) — ждём
// следующий тик, задачу не трогаем.
continue
}
w.captureInfohashes(ctx, d, t)
w.captureSourceAddedAt(ctx, d, t)
w.reconcile(ctx, d, t)
}
// Сверка разложенных задач с реальностью (источник в qBit + хардлинки на ФС)
// — отдельно от активных, по двумерной матрице (см. state-reconciliation).
w.reconcileDesync(ctx, byHash)
// Восстановление задач, упавших по нашей нетерпеливости (magnet_timeout/
// stalled), если их источник в qBittorrent ожил и продвинулся.
w.reconcileRecovery(ctx, byHash)
return nil
}
// reconcile двигает одну задачу по состоянию её торрента. Вызывается под
// w.mu.
func (w *Worker) reconcile(ctx context.Context, d store.Download, t qbt.Torrent) {
ctx = w.scoped(ctx, capIngest, d.ID, d.PrimaryInfohash())
switch classify(t.State) {
case classReady:
w.transition(ctx, d, store.StateCompleted, "", "")
case classErrored:
w.transition(ctx, d, store.StateFailed, errCodeQbitError, "qBittorrent state: "+t.State)
case classDownloading:
w.checkTimeouts(ctx, d, t)
case classBusy:
// moving/checking — ждём, файлы ещё не на финальном месте.
}
}
// checkTimeouts помечает зависшие задачи двумя разными мерами:
// - magnet_timeout — по ВОЗРАСТУ торрента (metaDL дольше magnet_timeout без
// метаданных): страховочный предохранитель (дефолт 24h), базис — added_on;
// - stuck_after — по ДЛИТЕЛЬНОСТИ ПРОСТОЯ (stalledDL без движения данных
// дольше stuck_after): базис — last_activity, а не возраст. Иначе долго
// качавшийся торрент, на миг зашедший в stalledDL, ложно уходит в stuck со
// «stalled for 5h» (см. state-reconciliation, MAJOR-2).
//
// Оба базиса приподняты до retried_at (ручной retry), чтобы возврат в
// downloading не ронял задачу снова на ближайшем тике. Настоящие провалы ловит
// classErrored, а ожившие задачи воскрешает reconcileRecovery.
func (w *Worker) checkTimeouts(ctx context.Context, d store.Download, t qbt.Torrent) {
switch {
case isMeta(t.State) && w.cfg.MagnetTimeout > 0:
if age, ok := w.torrentAge(d, t); ok && age > w.cfg.MagnetTimeout {
w.transition(ctx, d, store.StateFailed, errCodeMagnetTimeout,
fmt.Sprintf("no metadata after %s", age.Truncate(time.Second)))
}
case isStalledDL(t.State) && w.cfg.StuckAfter > 0:
if idle, ok := w.stallDuration(d, t); ok && idle > w.cfg.StuckAfter {
w.transition(ctx, d, store.StateStuck, errCodeStalled,
fmt.Sprintf("stalled for %s", idle.Truncate(time.Second)))
}
}
}
// captureSourceAddedAt однократно сохраняет время добавления торрента в
// qBittorrent (added_on) у задачи — базис сортировки списка. Пишем только при
// первом наблюдении (в БД source_added_at ещё пуст, SQL-гард в store); значение
// неизменно, поэтому повторные тики его не трогают. Учётная операция: её сбой не
// двигает задачу, лишь логируем WARN. Вызывается под w.mu.
func (w *Worker) captureSourceAddedAt(ctx context.Context, d store.Download, t qbt.Torrent) {
if d.SourceAddedAt.Valid || t.AddedOn <= 0 {
return
}
if err := w.store.SetSourceAddedAt(ctx, d.ID, time.Unix(t.AddedOn, 0)); err != nil {
w.log.Warn("capture source_added_at failed",
"capability", capIngest, "download_id", d.ID, "error", err)
}
}
// torrentFor ищет торрент загрузки в карте byHash по любому из её хешей.
func torrentFor(d store.Download, byHash map[string]qbt.Torrent) (qbt.Torrent, bool) {
for _, h := range d.HashList() {
if t, ok := byHash[h]; ok {
return t, true
}
}
return qbt.Torrent{}, false
}
// captureInfohashes дописывает загрузке хеши, которые qBittorrent знает, а мы
// ещё нет (гибридный торрент раскрывает v1+v2 после получения метаданных).
// Хеши собирает torrentHashes (усечённый t.Hash v2-only раздач отсеян).
// AddInfohashes под гардом: хеш, которым владеет другая активная задача,
// дописан не будет (ErrInfohashTaken). Учётная операция: сбой не двигает
// задачу, лишь логируем WARN. Под w.mu.
func (w *Worker) captureInfohashes(ctx context.Context, d store.Download, t qbt.Torrent) {
known := d.HashList()
var missing []string
for _, h := range torrentHashes(t) {
if !slices.Contains(known, h) {
missing = append(missing, h)
}
}
if len(missing) == 0 {
return
}
if err := w.store.AddInfohashes(ctx, d.ID, missing); err != nil {
w.log.Warn("capture infohashes failed",
"capability", capIngest, "download_id", d.ID, "error", err)
}
}
// torrentAge — возраст торрента для magnet_timeout: от added_on в qBittorrent
// (надёжный базис, переживает усыновление), с фолбэком на created_at задачи,
// если qBit не отдал added_on (NIT-10). Базис приподнят до retried_at, чтобы
// ручной retry сбрасывал отсчёт. ok=false — базис неизвестен (ни added_on, ни
// разбираемого created_at): таймаут не срабатывает, фиксируем диагностикой.
func (w *Worker) torrentAge(d store.Download, t qbt.Torrent) (time.Duration, bool) {
basis, ok := w.addedBasis(d, t)
if !ok {
return 0, false
}
return w.now().Sub(w.retriedFloor(d, basis)), true
}
// stallDuration — длительность простоя торрента для stuck_after: от
// last_activity qBittorrent (момент последнего движения данных), с фолбэком на
// базис добавления, если qBit не отдал пригодного last_activity. Базис приподнят
// до retried_at (ручной retry даёт свежее окно). ok=false — базис неизвестен.
//
// last_activity в будущем (перекос часов, sentinel «никогда не был активен»)
// трактуем как непригодное значение и падаем на addedBasis: иначе простой вышел
// бы отрицательным и реально застрявший торрент никогда бы не пометился stuck.
func (w *Worker) stallDuration(d store.Download, t qbt.Torrent) (time.Duration, bool) {
var basis time.Time
if la := time.Unix(t.LastActivity, 0).UTC(); t.LastActivity > 0 && !la.After(w.now()) {
basis = la
} else {
var ok bool
if basis, ok = w.addedBasis(d, t); !ok {
return 0, false
}
}
return w.now().Sub(w.retriedFloor(d, basis)), true
}
// addedBasis — момент добавления торрента: added_on qBittorrent, иначе
// created_at задачи (NIT-10). ok=false — ни того, ни другого разобрать не
// удалось; фиксируем диагностикой.
func (w *Worker) addedBasis(d store.Download, t qbt.Torrent) (time.Time, bool) {
if t.AddedOn > 0 {
return time.Unix(t.AddedOn, 0).UTC(), true
}
created, err := d.CreatedTime()
if err != nil {
w.log.Warn("cannot determine torrent age",
"capability", capIngest, "download_id", d.ID,
"created_at", d.CreatedAt, "error", err)
return time.Time{}, false
}
return created, true
}
// retriedFloor приподнимает базис отсчёта таймаута до времени последнего ручного
// retry: после retry задача получает свежее окно и не падает повторно на
// ближайшем тике (см. state-reconciliation «Ручной повтор», MAJOR-1).
func (w *Worker) retriedFloor(d store.Download, basis time.Time) time.Time {
if r, ok := d.RetriedTime(); ok && r.After(basis) {
return r
}
return basis
}
// transition пишет новое состояние и логирует переход. Fire-and-forget обёртка
// над transitionErr: применяется там, где переход терминален для шага — за ним
// нет побочного эффекта, зависящего от факта записи claim (reconcile, таймауты,
// команды ревью, финальные переходы linkPlan, sweep). Ошибку записи гасит (её
// уже залогировал transitionErr).
func (w *Worker) transition(ctx context.Context, d store.Download, state store.State, code, msg string) {
_ = w.transitionErr(ctx, d, state, code, msg)
}
// transitionErr пишет новое состояние, шлёт пинги/скан, логирует переход и
// ВОЗВРАЩАЕТ ошибку записи. На claim-then-side-effect путях (ручное Apply,
// авто-раскладка в finishRecognition) провал claim перехода в `linking` ОБЯЗАН
// прервать выполнение ДО побочных эффектов (хардлинков): иначе ссылки лягут при
// незакоммиченном claim, а финальный переход из фактического (не `linking`)
// состояния граф отклонит — задача застрянет со stale-планом (MINOR-7).
func (w *Worker) transitionErr(ctx context.Context, d store.Download, state store.State, code, msg string) error {
// FromOr, а не From: если вызывающий не завёл scoped-логгер, падаем на
// w.log (настроенный), а не на slog.Default().
log := logctx.FromOr(ctx, w.log)
if err := w.store.SetDownloadState(ctx, d.ID, state, code, msg); err != nil {
log.Error("state transition failed", "from", d.State, "to", state, "error", err)
return fmt.Errorf("transition %s → %s: %w", d.State, state, err)
}
log.Info("state transition", "from", d.State, "to", state, "code", code)
// Пинги — неблокирующе и в отдельном контексте: вызов уходит в сеть, а
// мы под w.mu (Notify читает состояние уже после освобождения замка).
if w.notifier != nil {
switch state {
case store.StateReview:
go w.notifier.Notify(context.Background(), d.ID, EventReview)
case store.StateDone:
go w.notifier.Notify(context.Background(), d.ID, EventDone)
case store.StateOrphaned:
go w.notifier.Notify(context.Background(), d.ID, EventOrphaned)
case store.StateTargetMissing:
go w.notifier.Notify(context.Background(), d.ID, EventTargetMissing)
case store.StateFailed, store.StateStuck:
if w.shouldNotifyFail(d.ID) {
go w.notifier.Notify(context.Background(), d.ID, EventFailed)
}
}
}
// Наши библиотечные хардлинки изменились — просим Jellyfin пересканировать
// библиотеку, чтобы плеер не держал битые пути и быстрее подхватил новые
// файлы. Триггерят входы, где раскладка «улеглась»: done (ссылки разложены),
// reverted (Undo снял ссылки), deleted (Delete снял / сверка констатировала
// отсутствие). target_missing/orphaned — промежуточный рассинхрон, ждём
// relink/лечения, не сканируем. Неблокирующе и вне w.mu; недоступность
// Jellyfin не влияет на состояние задачи.
if w.scanner != nil && triggersScan(state) {
// Скан Jellyfin — неблокирующе и вне w.mu, в фоновом ctx со scoped-логгером
// (download_id для корреляции ext.*-записи клиента). Недоступность Jellyfin
// на задачу не влияет; ошибку вызова логирует сам клиент (ext.*), здесь гасим.
gctx := w.scoped(context.Background(), capFileLayout, d.ID, d.PrimaryInfohash())
go func() { _ = w.scanner.RefreshLibraries(gctx) }()
}
return nil
}
// triggersScan сообщает, стоит ли на входе в state дёргать пересканирование
// Jellyfin: наши библиотечные хардлинки только что изменились. Гейт по
// состоянию-цели в едином чекпоинте ловит и пользовательские Undo/Delete, и
// reconcile-производный deleted (инициатор роли не играет); target_missing/
// orphaned — промежуточный рассинхрон (ждём relink/лечения) — исключены.
func triggersScan(state store.State) bool {
switch state {
case store.StateDone, store.StateReverted, store.StateDeleted:
return true
default:
return false
}
}
// shouldNotifyFail дебаунсит повторные уведомления о падении одной задачи
// (мерцающий stalled-торрент: stuck↔downloading), чтобы не спамить. Вызывается
// под w.mu. НЕ сбрасываем запись при восстановлении — иначе дебаунс не гасил бы
// флаппинг.
func (w *Worker) shouldNotifyFail(id string) bool {
now := w.now()
if last, ok := w.failNotified[id]; ok && now.Sub(last) < failNotifyDebounce {
return false
}
w.failNotified[id] = now
// Лёгкая чистка устаревших записей, чтобы карта не росла без предела.
for k, t := range w.failNotified {
if now.Sub(t) >= failNotifyDebounce {
delete(w.failNotified, k)
}
}
return true
}
// logCmd — единый чокпоинт логирования исхода команды воркера
// (apply/cancel/retry/…), вызываемой транспортами. Конвенция (logging.md,
// раздел «Ошибки»): доменную ошибку логирует граница домена ровно один раз, а
// транспорты (HTTP/web/Telegram) — нет. Команды воркера и есть эта граница.
//
// Уровень — по адресату: штатный отказ по состоянию/наличию
// (ErrConflict/ErrNotReady/ErrNotFound) адресован пользователю, он уже получил
// ответ на поверхности — DEBUG; всё прочее (сбой БД/ФС/зависимости) адресовано
// команде — ERROR. Успех (err == nil) — молча. Ставится в defer при именованном
// возврате, поэтому видит финальную ошибку и scoped-логгер, накопленный в ctx.
func (w *Worker) logCmd(ctx context.Context, cmd, id string, err error) {
if err == nil {
return
}
log := logctx.FromOr(ctx, w.log)
switch {
case errors.Is(err, ErrConflict), errors.Is(err, ErrNotReady), errors.Is(err, ErrInvalidInput),
errors.Is(err, store.ErrNotFound), errors.Is(err, layout.ErrCollision):
log.Debug("command rejected", "command", cmd, "download_id", id, "error", err)
default:
log.Error("command failed", "command", cmd, "download_id", id, "error", err)
}
}
// Cancel отклоняет задачу. Торрент в qBittorrent не трогаем — он продолжает
// раздачу (источник неприкосновенен).
func (w *Worker) Cancel(ctx context.Context, id string) (err error) {
defer func() { w.logCmd(ctx, "cancel", id, err) }()
w.mu.Lock()
defer w.mu.Unlock()
d, err := w.store.GetDownload(ctx, id)
if err != nil {
return fmt.Errorf("cancel: %w", err)
}
if d.State.IsTerminal() {
return fmt.Errorf("cancel: download %s is already terminal (%s): %w", id, d.State, ErrConflict)
}
if err := w.store.SetDownloadState(ctx, id, store.StateCancelled, "", ""); err != nil {
return fmt.Errorf("cancel: %w", err)
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("state transition",
"from", d.State, "to", store.StateCancelled)
return nil
}
// Dismiss — универсальный стоп-кран: переводит задачу в терминальный cancelled из
// ЛЮБОГО состояния, кроме deleted, ТОЛЬКО меняя статус. В отличие от Delete не
// трогает ни файлы (библиотечные хардлинки done/orphaned остаются на месте), ни
// раздачу в qBittorrent, ни цель; source-preflight не делает. Служит закрытием
// зависшей/спорной/лишней записи (в т.ч. дубля-близнеца в target_missing). Из
// cancelled — идемпотентный no-op БЕЗ setState: иначе перезаписал бы error_code,
// подменив причину прежнего Cancel/Dismiss. Помечает переход user_dismiss
// (отличает от reconcile и от штатного Cancel с пустым кодом).
func (w *Worker) Dismiss(ctx context.Context, id string) (err error) {
defer func() { w.logCmd(ctx, "dismiss", id, err) }()
w.mu.Lock()
defer w.mu.Unlock()
d, err := w.store.GetDownload(ctx, id)
if err != nil {
return fmt.Errorf("dismiss: %w", err)
}
if d.State == store.StateDeleted {
return fmt.Errorf("dismiss: download %s is deleted (strictly terminal): %w", id, ErrConflict)
}
if d.State == store.StateCancelled {
return nil // уже закрыта — no-op, error_code прежней отмены не трогаем
}
if err := w.store.SetDownloadState(ctx, id, store.StateCancelled, "user_dismiss", "закрыто пользователем"); err != nil {
return fmt.Errorf("dismiss: %w", err)
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("state transition",
"from", d.State, "to", store.StateCancelled, "code", "user_dismiss")
return nil
}
// Retry повторяет застрявшую/упавшую задачу: заново отдаёт источник в
// qBittorrent и возвращает в downloading.
func (w *Worker) Retry(ctx context.Context, id string) (err error) {
defer func() { w.logCmd(ctx, "retry", id, err) }()
w.mu.Lock()
defer w.mu.Unlock()
d, err := w.store.GetDownload(ctx, id)
if err != nil {
return fmt.Errorf("retry: %w", err)
}
if d.State != store.StateFailed && d.State != store.StateStuck {
return fmt.Errorf("retry: download %s is %s, only failed/stuck are retriable: %w", id, d.State, ErrConflict)
}
// Если раздача уже жива и ЗДОРОВА в qBittorrent — перецепляемся к ней,
// повторный Add не нужен (и вреден: вслепую дублировал бы торрент). Add — когда
// источника в qBittorrent нет. Базис таймаута сбрасывается ниже через
// retried_at, поэтому возврат в downloading не роняет задачу снова на ближайшем
// тике (MAJOR-1).
reAdd := true
if hashes := d.HashList(); len(hashes) > 0 {
var t qbt.Torrent
var alive bool
t, alive, err = w.torrentByInfohash(ctx, hashes)
if err != nil {
return fmt.Errorf("retry: %w", err)
}
if alive && classify(t.State) == classErrored {
// Живой, но сломанный торрент (error/missingFiles): повторный Add его не
// чинит (qBittorrent отверг бы дубль), а простой возврат в downloading
// тут же снова упал бы classErrored на ближайшем тике сверки (+ дебаунс
// уведомления) — retry выглядел бы сломанным. Отклоняем с подсказкой
// починить раздачу в qBittorrent (recheck/восстановить файлы) — NIT-12.
return fmt.Errorf("retry: раздача в qBittorrent в состоянии ошибки (%s) — почини её (recheck) в qBittorrent и повтори: %w", t.State, ErrConflict)
}
reAdd = !alive
}
// Гард инварианта — ДО побочного эффекта в qBittorrent: пока задача лежала
// в failed, тем же infohash могла завладеть другая активная задача — тогда
// отказываем, не добавив торрент повторно (см. design ulid-identity, D4).
if err := w.store.ActivateIfNoOtherActive(ctx, id, store.StateDownloading, "", ""); err != nil {
if errors.Is(err, store.ErrInfohashTaken) {
return fmt.Errorf("retry: для этого торрента уже есть другая активная задача: %w", ErrConflict)
}
return fmt.Errorf("retry: %w", err)
}
if reAdd {
// Добавляем заново по типу источника (magnet — ссылкой, torrent —
// сохранёнными байтами файлом). Rename при retry не выводим (namer здесь
// не зовём — имя уже могло быть выведено при первом добавлении).
_, addReq, prepErr := w.sourceAddParts(ctx, *d)
if prepErr != nil {
// Байты torrent недоступны — откатываем активацию, задача не должна
// «качаться» без раздачи в qBittorrent.
if rbErr := w.store.SetDownloadState(ctx, id, d.State, d.ErrorCode.String, d.ErrorMsg.String); rbErr != nil {
w.log.Error("retry rollback failed",
"capability", capReview, "download_id", id, "error", rbErr)
}
return fmt.Errorf("retry: prepare add: %w", prepErr)
}
if err := w.qbt.Add(ctx, addReq); err != nil {
// Активация уже прошла — откатываем задачу в прежнее состояние,
// чтобы не оставить «качающуюся» задачу без раздачи в qBittorrent.
if rbErr := w.store.SetDownloadState(ctx, id, d.State, d.ErrorCode.String, d.ErrorMsg.String); rbErr != nil {
w.log.Error("retry rollback failed",
"capability", capReview, "download_id", id, "error", rbErr)
}
return fmt.Errorf("retry: add to qbittorrent: %w", err)
}
}
// Сброс базиса отсчёта таймаутов (MAJOR-1): без него живой, но давно
// добавленный/простаивающий торрент снова упал бы по magnet_timeout/
// stuck_after на ближайшем тике. Best-effort: сбой лишь лишает свежего окна
// (WARN), сам retry уже состоялся.
if err := w.store.SetRetriedAt(ctx, id, w.now()); err != nil {
w.log.Warn("retry basis reset failed",
"capability", capReview, "download_id", id, "error", err)
}
// Сброс счётчика пропусков источника: retried source_gone-задача иначе вошла бы
// в downloading с source_miss_count == threshold и упала бы снова на ближайшем
// тике, если переотданная раздача ещё не видна в выдаче qBittorrent — без
// обещанного грейс-окна (MAJOR-3). Best-effort: сбой лишь лишает свежего окна.
if d.SourceMissCount != 0 {
if err := w.store.SetSourceMissCount(ctx, id, 0); err != nil {
w.log.Warn("retry miss count reset failed",
"capability", capReview, "download_id", id, "error", err)
}
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("state transition",
"from", d.State, "to", store.StateDownloading)
return nil
}
// class — класс состояния торрента qBittorrent.
type class int
const (
classDownloading class = iota // ещё качается
classReady // готов к раскладке
classErrored // ошибка
classBusy // moving/checking — переходный момент, ждём
)
// classify относит состояние qBittorrent к классу (см. architecture.md,
// «Завершение в qBittorrent»). Учитываем и v5-имена (stopped* вместо
// paused*).
func classify(state string) class {
switch state {
case "uploading", "stalledUP", "pausedUP", "stoppedUP", "queuedUP", "forcedUP":
return classReady
case "error", "missingFiles":
return classErrored
case "moving", "checkingUP", "checkingResumeData", "allocating":
return classBusy
default:
// downloading, stalledDL, metaDL, forcedMetaDL, queuedDL, checkingDL,
// forcedDL, pausedDL, stoppedDL, unknown — считаем «ещё качается».
return classDownloading
}
}
func isMeta(state string) bool {
return state == "metaDL" || state == "forcedMetaDL"
}
func isStalledDL(state string) bool {
return state == "stalledDL"
}