Долгий metaDL больше не убивается агрессивным таймаутом: дефолт magnet_timeout 30m → 24h (страховочный предохранитель), базис отсчёта — added_on из qBittorrent, а не created_at (переживает retry/усыновление). Авто-восстановление: фоновая сверка возвращает в поток задачи, упавшие по нашей нетерпеливости (magnet_timeout/stalled), когда источник ожил и продвинулся за условие падения (downloading/completed по статусу торрента); qbit_error не воскрешается. Конфликт idempotency (infohash занят другой активной задачей) — оставляем в failed. Уведомления: любой переход в failed/stuck пингует автора (включая приёмный qbit_add через ingest), с дебаунсом против спама при флаппинге stalled. Ручной retry добавлен в веб-UI и Telegram; Retry перецепляется к живому торренту вместо слепого Add. Дельта state-reconciliation влита в живые спеки; обновлён workflow.md. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
289 lines
13 KiB
Go
289 lines
13 KiB
Go
package worker
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
|
|
"git.vakhrushev.me/av/jellybit/internal/layout"
|
|
"git.vakhrushev.me/av/jellybit/internal/logctx"
|
|
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
|
"git.vakhrushev.me/av/jellybit/internal/store"
|
|
)
|
|
|
|
// desyncStates — состояния, которые ведёт сверка с реальностью: уже
|
|
// разложенные (done) и восстановимые состояния рассинхрона. deleted сюда не
|
|
// входит — оно терминально и сверкой не переоценивается (источник к
|
|
// терминальной задаче не вернётся из-за идемпотентности, а цель отбирается
|
|
// переходом владения путём; см. state-reconciliation). Активные и
|
|
// пользовательски-терминальные (reverted/cancelled/failed/stuck) сверка не
|
|
// трогает.
|
|
var desyncStates = []store.State{
|
|
store.StateDone,
|
|
store.StateTargetMissing,
|
|
store.StateOrphaned,
|
|
}
|
|
|
|
// deriveState выводит состояние задачи из двумерной матрицы «источник × цель»
|
|
// (см. state-reconciliation, D1).
|
|
func deriveState(sourcePresent, targetPresent bool) store.State {
|
|
switch {
|
|
case sourcePresent && targetPresent:
|
|
return store.StateDone
|
|
case sourcePresent && !targetPresent:
|
|
return store.StateTargetMissing
|
|
case !sourcePresent && targetPresent:
|
|
return store.StateOrphaned
|
|
default:
|
|
return store.StateDeleted
|
|
}
|
|
}
|
|
|
|
// reconcileDesync сверяет разложенные/desync-задачи с реальностью. Вызывается
|
|
// из Poll под w.mu. byHash — карта присутствующих в qBittorrent раздач.
|
|
func (w *Worker) reconcileDesync(ctx context.Context, byHash map[string]qbt.Torrent) {
|
|
cands, err := w.store.ListDownloadsByState(ctx, desyncStates...)
|
|
if err != nil {
|
|
w.log.Warn("reconcile list desync failed", "capability", capIngest, "error", err)
|
|
return
|
|
}
|
|
for _, d := range cands {
|
|
w.reconcileOneDesync(ctx, d, byHash)
|
|
}
|
|
}
|
|
|
|
// reconcileOneDesync сверяет одну задачу: вычисляет присутствие источника (с
|
|
// дебаунсом) и цели, выводит состояние и переходит при изменении.
|
|
func (w *Worker) reconcileOneDesync(ctx context.Context, d store.Download, byHash map[string]qbt.Torrent) {
|
|
if !d.Infohash.Valid {
|
|
return // нечем сопоставить источник
|
|
}
|
|
ctx = w.scoped(ctx, capIngest, d.ID, d.Infohash.String)
|
|
_, sourceSeen := byHash[strings.ToLower(d.Infohash.String)]
|
|
|
|
// Дебаунс пропажи источника: считаем удалённым только после порога подряд
|
|
// идущих промахов; любое появление сбрасывает счётчик.
|
|
sourcePresent := w.debounceSource(ctx, d, sourceSeen)
|
|
|
|
targetPresent, err := w.targetPresent(ctx, d.ID)
|
|
if err != nil {
|
|
logctx.From(ctx).Warn("reconcile target probe failed", "error", err)
|
|
return // не можем проверить цель — не дёргаем состояние
|
|
}
|
|
|
|
want := deriveState(sourcePresent, targetPresent)
|
|
if want == d.State {
|
|
return
|
|
}
|
|
w.transition(ctx, d, want, "reconcile", reconcileReason(sourcePresent, targetPresent))
|
|
}
|
|
|
|
// debounceSource обновляет счётчик промахов источника и возвращает, считать ли
|
|
// источник присутствующим с учётом порога. Запись счётчика — только при его
|
|
// изменении (без лишних UPDATE на каждом тике).
|
|
func (w *Worker) debounceSource(ctx context.Context, d store.Download, sourceSeen bool) bool {
|
|
threshold := max(w.cfg.SourceMissingThreshold, 1)
|
|
if sourceSeen {
|
|
if d.SourceMissCount != 0 {
|
|
if err := w.store.SetSourceMissCount(ctx, d.ID, 0); err != nil {
|
|
logctx.From(ctx).Warn("reconcile reset miss count failed", "error", err)
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
miss := d.SourceMissCount + 1
|
|
if err := w.store.SetSourceMissCount(ctx, d.ID, miss); err != nil {
|
|
logctx.From(ctx).Warn("reconcile bump miss count failed", "error", err)
|
|
}
|
|
// До порога источник трактуется как присутствующий — задача не дёргается.
|
|
return miss < threshold
|
|
}
|
|
|
|
// targetPresent сообщает, существуют ли разложенные хардлинки задачи. Цель
|
|
// считается присутствующей, только если существуют ВСЕ ссылки последнего
|
|
// батча; частичная пропажа — это отсутствие цели (библиотека сломана → relink).
|
|
func (w *Worker) targetPresent(ctx context.Context, id int64) (bool, error) {
|
|
batch, err := w.store.LatestBatchID(ctx, id)
|
|
if err != nil {
|
|
return false, fmt.Errorf("latest batch: %w", err)
|
|
}
|
|
if batch == "" {
|
|
return false, nil // ничего не разложено — цели нет
|
|
}
|
|
rows, err := w.store.ListFileLinksByBatch(ctx, batch)
|
|
if err != nil {
|
|
return false, fmt.Errorf("list links: %w", err)
|
|
}
|
|
laidOut := 0
|
|
for _, r := range rows {
|
|
// Все статусы, означающие реальный файл на диске: только что слинкован,
|
|
// скопирован (фолбэк) или уже существовал тем же inode (идемпотентный
|
|
// повтор apply). Прочие (collision и т.п.) — не наша разложенная цель.
|
|
if !isLaidOut(r.Status) {
|
|
continue
|
|
}
|
|
laidOut++
|
|
if _, err := os.Lstat(r.DstPath); err != nil {
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return false, nil
|
|
}
|
|
return false, fmt.Errorf("stat target %q: %w", r.DstPath, err)
|
|
}
|
|
}
|
|
return laidOut > 0, nil
|
|
}
|
|
|
|
// isLaidOut сообщает, означает ли статус file_link реально лежащий на ФС файл
|
|
// нашей раскладки (а не коллизию или иной не-разложенный исход).
|
|
func isLaidOut(status string) bool {
|
|
switch layout.LinkStatus(status) {
|
|
case layout.StatusLinked, layout.StatusCopied, layout.StatusExists:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// reconcileReason — человекочитаемая причина перехода для лога/UI.
|
|
func reconcileReason(sourcePresent, targetPresent bool) string {
|
|
switch {
|
|
case sourcePresent && targetPresent:
|
|
return "источник и цель на месте"
|
|
case sourcePresent && !targetPresent:
|
|
return "цель удалена, источник на месте"
|
|
case !sourcePresent && targetPresent:
|
|
return "источник удалён, цель (последняя копия) на месте"
|
|
default:
|
|
return "удалены и источник, и цель"
|
|
}
|
|
}
|
|
|
|
// --- Восстановление зависших загрузок (failed/stuck → поток) ---
|
|
|
|
// reconcileRecovery воскрешает задачи, упавшие из-за нашей нетерпеливости
|
|
// (magnet_timeout/stalled), когда их источник в qBittorrent ожил и продвинулся
|
|
// за условие падения. Вызывается из Poll под w.mu. Реальные/пользовательские
|
|
// провалы (qbit_error/reverted/cancelled/deleted) сюда не попадают.
|
|
func (w *Worker) reconcileRecovery(ctx context.Context, byHash map[string]qbt.Torrent) {
|
|
cands, err := w.store.ListRecoverable(ctx, errCodeMagnetTimeout, errCodeStalled)
|
|
if err != nil {
|
|
w.log.Warn("recovery list failed", "capability", capIngest, "error", err)
|
|
return
|
|
}
|
|
for _, d := range cands {
|
|
w.reconcileOneRecovery(ctx, d, byHash)
|
|
}
|
|
}
|
|
|
|
// reconcileOneRecovery возвращает одну зависшую задачу в поток, если её торрент
|
|
// присутствует и продвинулся за условие падения.
|
|
func (w *Worker) reconcileOneRecovery(ctx context.Context, d store.Download, byHash map[string]qbt.Torrent) {
|
|
if !d.Infohash.Valid {
|
|
return
|
|
}
|
|
t, ok := byHash[strings.ToLower(d.Infohash.String)]
|
|
if !ok {
|
|
return // источника нет — оставляем как есть (вернёт ручной retry)
|
|
}
|
|
if !torrentProgressed(d, t) {
|
|
return // торрент всё ещё в metaDL/stalledDL или ошибочен — не воскрешаем
|
|
}
|
|
want := recoveredState(t.State)
|
|
if want == "" {
|
|
return // переходное состояние qBit (moving/checking) — ждём
|
|
}
|
|
ctx = w.scoped(ctx, capIngest, d.ID, d.Infohash.String)
|
|
|
|
// Конфликт идемпотентности: пока задача лежала в failed, тот же infohash мог
|
|
// взять другая активная задача (idempotency_key снят при падении). Оба
|
|
// целевых состояния (downloading/completed) нетерминальны → SetDownloadState
|
|
// восстановит idempotency_key = infohash; при занятом ключе упёрлись бы в
|
|
// unique-индекс. Поэтому проверяем владельца независимо от целевого состояния
|
|
// и оставляем старую задачу в failed.
|
|
other, err := w.store.FindActiveByInfohash(ctx, d.Infohash.String)
|
|
if err != nil {
|
|
logctx.From(ctx).Warn("recovery active lookup failed", "error", err)
|
|
return
|
|
}
|
|
if other != nil && other.ID != d.ID {
|
|
logctx.From(ctx).Info("recovery skipped, infohash taken by active download",
|
|
"conflict_download_id", other.ID)
|
|
return
|
|
}
|
|
// error_code/error_msg не пишем — задача снова здорова; причину в лог, а не в
|
|
// поле ошибки (иначе она светилась бы в UI/REST как ошибка живой задачи).
|
|
logctx.From(ctx).Info("recovery from failure", "to", want, "qbit_state", t.State)
|
|
w.transition(ctx, d, want, "", "")
|
|
}
|
|
|
|
// torrentProgressed сообщает, продвинулся ли торрент за условие, по которому
|
|
// задача упала: для magnet_timeout — получил метаданные (вышел из metaDL); для
|
|
// stalled — раздача ожила (вышла из stalledDL). Ошибочные состояния qBittorrent
|
|
// продвижением не считаем (их ведёт обычный reconcile в qbit_error).
|
|
func torrentProgressed(d store.Download, t qbt.Torrent) bool {
|
|
if classify(t.State) == classErrored {
|
|
return false
|
|
}
|
|
switch d.ErrorCode.String {
|
|
case errCodeMagnetTimeout:
|
|
return !isMeta(t.State)
|
|
case errCodeStalled:
|
|
return !isStalledDL(t.State)
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// recoveredState выводит состояние воскрешённой задачи из состояния торрента:
|
|
// готов к раскладке → completed; ещё качается → downloading. Переходные
|
|
// (moving/checking) и ошибочные состояния не восстанавливаем (пусто).
|
|
func recoveredState(state string) store.State {
|
|
switch classify(state) {
|
|
case classReady:
|
|
return store.StateCompleted
|
|
case classDownloading:
|
|
return store.StateDownloading
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
// --- Синхронный preflight перед действием (не доверяем state в БД) ---
|
|
|
|
// ensureSourcePresent синхронно (без дебаунса) проверяет, что раздача есть в
|
|
// qBittorrent прямо сейчас. При отсутствии приводит состояние к реальности и
|
|
// возвращает ErrConflict. Недоступность qBittorrent — честный отказ операции.
|
|
func (w *Worker) ensureSourcePresent(ctx context.Context, d *store.Download, op string) error {
|
|
if !d.Infohash.Valid {
|
|
return fmt.Errorf("%s: download %d has no infohash", op, d.ID)
|
|
}
|
|
_, ok, err := w.torrentByInfohash(ctx, d.Infohash.String)
|
|
if err != nil {
|
|
return fmt.Errorf("%s: %w", op, err)
|
|
}
|
|
if ok {
|
|
return nil
|
|
}
|
|
// Источник пропал — немедленно приводим состояние к реальности.
|
|
w.reconcileToReality(ctx, *d, false)
|
|
return fmt.Errorf("%s: источник удалён из qBittorrent: %w", op, ErrConflict)
|
|
}
|
|
|
|
// reconcileToReality выводит и проставляет состояние по уже известному факту об
|
|
// источнике (sourcePresent) и фактически проверенной цели. Используется
|
|
// preflight-проверками: немедленно, без дебаунса.
|
|
func (w *Worker) reconcileToReality(ctx context.Context, d store.Download, sourcePresent bool) {
|
|
targetPresent, err := w.targetPresent(ctx, d.ID)
|
|
if err != nil {
|
|
logctx.FromOr(ctx, w.log).Warn("preflight target probe failed", "error", err)
|
|
return
|
|
}
|
|
want := deriveState(sourcePresent, targetPresent)
|
|
if want == d.State {
|
|
return
|
|
}
|
|
w.transition(ctx, d, want, "reconcile", reconcileReason(sourcePresent, targetPresent))
|
|
}
|