Отмена задачи (catched→cancelled) в окно, пока worker вне блокировки выводит имя (LLM) и делает qbt.Add, оставляла добавленный торрент в qBittorrent без задачи-владельца: PromoteCatched корректно пропускал переход, но источник уже качался/сидировал вечно, а усыновить его назад нельзя (хеши принадлежат отменённой задаче). Спека покрывала переход состояния, но не побочный эффект. Комбинированная защита в processCatched: - re-read состояния под w.mu прямо перед qbt.Add — при отмене источник не добавляется вовсе (сужает окно гонки); - свежий листинг перед add подтверждает отсутствие infohash — признак «своего» торрента; при сбое листинга/присутствии add не делаем (усыновит следующий тик); - при отмене в окне после add (промах PromoteCatched, подтверждённый re-read'ом state != catched) — уборка добавленного нами торрента qbt.Delete(_, true); - WARN/ERROR-логи по этому пути с корреляцией по download_id, без секретов. Гарантия «удаляем только своё»: удаление-с-данными достижимо ТОЛЬКО после подтверждённого отсутствия infohash перед add, поэтому пред-существующий/чужой торрент с тем же хешем никогда не сносится (негативный инвариант). Обоснование по инварианту «источник неприкосновенен» — в design.md изменения. Дельта — download-tracking (требование «Добавление пойманной загрузки в qBittorrent»): re-read перед add, подтверждение отсутствия, уборка при отмене, негативный сценарий. Тесты покрывают все ветки (skip-before-add, cleanup после add, пред-существующий не удаляется, сбой БД не удаляет, сбой листинга не добавляет). Change archived: 2026-07-17-cancel-during-add-cleanup. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1134 lines
62 KiB
Go
1134 lines
62 KiB
Go
// 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 пересканировать библиотеку, чтобы
|
||
// новые файлы быстрее появились в проигрывателе. Тоже неблокирующе и вне
|
||
// w.mu; недоступность Jellyfin не влияет на состояние задачи.
|
||
if w.scanner != nil && state == store.StateDone {
|
||
// Скан 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
|
||
}
|
||
|
||
// 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"
|
||
}
|