package worker import ( "context" "errors" "fmt" "os" "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 len(d.Infohashes) == 0 { return // нечем сопоставить источник } ctx = w.scoped(ctx, capIngest, d.ID, d.PrimaryInfohash()) _, sourceSeen := torrentFor(d, byHash) // Дебаунс пропажи источника: считаем удалённым только после порога подряд // идущих промахов; любое появление сбрасывает счётчик. 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 string) (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 len(d.Infohashes) == 0 { return } t, ok := torrentFor(d, byHash) 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.PrimaryInfohash()) // Возврат в активное состояние — только через атомарный гард инварианта: // пока задача лежала в failed, тем же infohash могла завладеть другая // активная задача (новый приём). Тогда старую оставляем в failed. // error_code/error_msg не пишем — задача снова здорова; причину в лог, а не в // поле ошибки (иначе она светилась бы в UI/REST как ошибка живой задачи). if err := w.store.ActivateIfNoOtherActive(ctx, d.ID, want, "", ""); err != nil { if errors.Is(err, store.ErrInfohashTaken) { logctx.From(ctx).Info("recovery skipped, infohash taken by active download", "error", err) return } logctx.From(ctx).Warn("recovery activate failed", "error", err) return } logctx.From(ctx).Info("recovery from failure", "from", d.State, "to", want, "qbit_state", t.State) } // 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 в БД) --- // ensureSourceReady синхронно (без дебаунса) проверяет, что раздача есть в // qBittorrent прямо сейчас И докачана (класс classReady). Команды ревью, ведущие // к распознаванию или раскладке, работают только с готовым источником — иначе // хардлинки легли бы на неполные файлы (qBittorrent отдаёт имена до завершения). // - источник исчез → приводим состояние к реальности (orphaned/deleted) и // ErrConflict; // - источник есть, но ещё качается → ErrNotReady БЕЗ reconcile: нахождение // задачи в review/deferred/… легитимно, приводить нечего. // // Недоступность qBittorrent — честный отказ операции. func (w *Worker) ensureSourceReady(ctx context.Context, d *store.Download, op string) error { if len(d.Infohashes) == 0 { return fmt.Errorf("%s: download %s has no infohash", op, d.ID) } t, ok, err := w.torrentByInfohash(ctx, d.HashList()) if err != nil { return fmt.Errorf("%s: %w", op, err) } if !ok { // Источник пропал — немедленно приводим состояние к реальности. w.reconcileToReality(ctx, *d, false) return fmt.Errorf("%s: источник удалён из qBittorrent: %w", op, ErrConflict) } if classify(t.State) != classReady { return fmt.Errorf("%s: торрент ещё качается: %w", op, ErrNotReady) } return nil } // 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)) }