// 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 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 } // Успех: короткий переход под w.mu с ре-валидацией state=catched // (загрузку могли отменить, пока шли сетевые вызовы). w.mu.Lock() if err := w.store.PromoteCatched(cctx, d.ID, rename); err != nil { logctx.From(cctx).Info("catched promote skipped", "reason", err.Error()) } else { logctx.From(cctx).Info("state transition", "from", store.StateCatched, "to", store.StateDownloading) } w.mu.Unlock() } } // 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 []string{t.Hash, t.InfohashV1, t.InfohashV2} { if h != "" { 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 []string{t.Hash, t.InfohashV1, t.InfohashV2} { if h != "" { 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 нет ИЛИ он в состоянии ошибки: перецепка к // сломанному торренту (error/missingFiles) бессмысленна — reconcile тут же // вернул бы задачу в failed, поэтому пробуем повторно отдать источник // (NIT-12). Базис таймаута сбрасывается ниже через 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) } reAdd = !alive || classify(t.State) == classErrored } // Гард инварианта — ДО побочного эффекта в 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" }