На /download/{id} поле «Режиссёр» было захардкожено прочерком, хотя экран
ревью режиссёра уже выводит: слоистое разрешение полей ярлыка было заперто в
неэкспортируемом worker.effectiveDisplayName. Из-за этого билдеры вью видели
только слой распознавания+матч (rd.Plan.Director) без слоя контекста — то же
на экране ревью.
Вынес разрешение в экспортируемую naming.EffectiveFields(parsedContext, plan)
LabelFields с методом Label(): выбор слоя по сырым значениям (как прежде),
выбранные скаляры возвращаются очищенными (sanitize идемпотентен, display_name
побайтно тот же). effectiveDisplayName стал тонкой обёрткой; страница загрузки
и экран ревью берут режиссёра из той же функции — согласованно с заголовком.
OpenSpec: web-ui (ADDED «Режиссёр в блоке распознавания страницы загрузки»),
review (MODIFIED «Инфо и предпросмотр выбранного источника» — слоистое
разрешение с фолбэком на контекст). Change заархивирован. Беклог: закрыта
rezhisser-v-kartochke-zagruzki.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1442 lines
65 KiB
Go
1442 lines
65 KiB
Go
package worker
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"maps"
|
||
"os"
|
||
"path/filepath"
|
||
"strconv"
|
||
"strings"
|
||
|
||
"git.vakhrushev.me/av/jellybit/internal/layout"
|
||
"git.vakhrushev.me/av/jellybit/internal/logctx"
|
||
"git.vakhrushev.me/av/jellybit/internal/metadata"
|
||
"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"
|
||
)
|
||
|
||
// Поля override.
|
||
const (
|
||
ovrMediaType = "media_type"
|
||
ovrIgnoredFiles = "ignored_files"
|
||
ovrProvider = "provider" // выбранная база ("none" = без базы)
|
||
ovrProviderID = "provider_id" // id в выбранной базе
|
||
ovrTitle = "title" // запиненное каноническое название
|
||
ovrYear = "year" // запиненный год
|
||
ovrDirector = "director" // режиссёр выбранного источника (credits метабазы)
|
||
ovrForceReview = "force_review" // ручная перепривязка: не авто-раскладывать
|
||
)
|
||
|
||
// recognizePending распознаёт завершённые загрузки и перезапускает те, что
|
||
// помечены к перераспознаванию (recognizing — например, после подсказки или
|
||
// после рестарта сервиса). Выполняется последовательно в поллинг-горутине;
|
||
// сам вызов LLM идёт вне блокировки, поэтому команды ревью не простаивают.
|
||
func (w *Worker) recognizePending(ctx context.Context) {
|
||
w.mu.Lock()
|
||
pending, err := w.store.ListDownloadsByState(ctx, store.StateCompleted, store.StateRecognizing)
|
||
w.mu.Unlock()
|
||
if err != nil {
|
||
w.log.Warn("recognition list pending failed", "error", err)
|
||
return
|
||
}
|
||
for _, d := range pending {
|
||
w.recognizeOne(ctx, d.ID)
|
||
}
|
||
}
|
||
|
||
// recognizeOne проводит одну загрузку через распознавание. Claim-паттерн:
|
||
// под блокировкой переводим в recognizing, LLM зовём без блокировки, затем
|
||
// под блокировкой фиксируем результат — но только если задачу за это время
|
||
// не увели в другое состояние (cancel/defer).
|
||
func (w *Worker) recognizeOne(ctx context.Context, id string) {
|
||
w.mu.Lock()
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
w.mu.Unlock()
|
||
w.log.Warn("recognition get download failed", "download_id", id, "error", err)
|
||
return
|
||
}
|
||
if d.State != store.StateCompleted && d.State != store.StateRecognizing {
|
||
w.mu.Unlock()
|
||
return
|
||
}
|
||
ctx = w.scoped(ctx, capRecognize, id, d.PrimaryInfohash())
|
||
if d.State == store.StateCompleted {
|
||
w.transition(ctx, *d, store.StateRecognizing, "", "")
|
||
// Перечитываем ради свежего updated_at: он служит claim-токеном этой
|
||
// попытки (MINOR-8). Токен фиксирует конкретный recognizing-эпизод; если
|
||
// задачу позже уведут из recognizing и вернут обратно (cancel → relink),
|
||
// updated_at сменится, и finishRecognition отбросит устаревший результат.
|
||
d, err = w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
w.mu.Unlock()
|
||
logctx.From(ctx).Warn("recognition reload after claim failed", "error", err)
|
||
return
|
||
}
|
||
}
|
||
claim := d.UpdatedAt
|
||
w.mu.Unlock()
|
||
|
||
result, savePath, err := w.runRecognize(ctx, *d)
|
||
if err != nil {
|
||
// Граница доменной стадии распознавания: логируем исход один раз (ERROR),
|
||
// дальше уходим в review с причиной — человек перезапустит подсказкой.
|
||
logctx.From(ctx).Error("recognition failed", "error", err)
|
||
result = recognize.Result{Decision: recognize.Decision{
|
||
Reasons: []string{"распознавание не удалось: " + err.Error()},
|
||
}}
|
||
}
|
||
w.finishRecognition(ctx, id, claim, result, savePath)
|
||
}
|
||
|
||
// runRecognize собирает сигналы из qBittorrent и накопленные подсказки,
|
||
// затем зовёт распознаватель. Возвращает также savePath для маппинга
|
||
// относительных путей файлов в абсолютные при раскладке.
|
||
func (w *Worker) runRecognize(ctx context.Context, d store.Download) (recognize.Result, string, error) {
|
||
if len(d.Infohashes) == 0 {
|
||
return recognize.Result{}, "", fmt.Errorf("no infohash")
|
||
}
|
||
t, ok, err := w.torrentByInfohash(ctx, d.HashList())
|
||
if err != nil {
|
||
return recognize.Result{}, "", err
|
||
}
|
||
if !ok {
|
||
return recognize.Result{}, "", fmt.Errorf("torrent not found in qBittorrent")
|
||
}
|
||
files, err := w.qbt.Files(ctx, t.Hash)
|
||
if err != nil {
|
||
return recognize.Result{}, "", err
|
||
}
|
||
hints, err := w.store.ListHints(ctx, d.ID)
|
||
if err != nil {
|
||
return recognize.Result{}, "", err
|
||
}
|
||
|
||
in := recognize.Input{
|
||
Name: t.Name,
|
||
Context: d.Context,
|
||
Hints: hints,
|
||
Files: make([]recognize.File, len(files)),
|
||
}
|
||
for i, f := range files {
|
||
in.Files[i] = recognize.File{Path: f.Name, Size: f.Size}
|
||
}
|
||
|
||
savePath := translatePath(t.SavePath, w.cfg.PathMap)
|
||
res, err := w.recognizer.Recognize(ctx, in)
|
||
if err != nil {
|
||
return recognize.Result{}, savePath, err
|
||
}
|
||
return res, savePath, nil
|
||
}
|
||
|
||
// finishRecognition сохраняет попытку распознавания и двигает задачу. В Ф3
|
||
// метабазы выключены → авто-раскладки не делаем, всегда уходим в review.
|
||
func (w *Worker) finishRecognition(ctx context.Context, id, claim string, res recognize.Result, savePath string) {
|
||
log := logctx.From(ctx)
|
||
planJSON, err := json.Marshal(res.Plan)
|
||
if err != nil {
|
||
log.Error("recognition marshal plan failed", "error", err)
|
||
planJSON = []byte("{}")
|
||
}
|
||
|
||
provider, providerID := "none", ""
|
||
if res.Match != nil {
|
||
provider, providerID = res.Match.Provider, res.Match.ProviderID
|
||
}
|
||
|
||
rec := &store.Recognition{
|
||
DownloadID: id,
|
||
MediaType: store.NullString(string(res.Plan.Type)),
|
||
Title: store.NullString(res.Plan.Title),
|
||
Provider: store.NullString(provider),
|
||
ProviderID: store.NullString(providerID),
|
||
Plan: store.NullString(string(planJSON)),
|
||
RawLLM: store.NullString(res.Raw),
|
||
}
|
||
if res.Plan.OriginalTitle != "" {
|
||
rec.OriginalTitle = store.NullString(res.Plan.OriginalTitle)
|
||
}
|
||
if res.Plan.Year != 0 {
|
||
rec.Year = sql.NullInt64{Int64: int64(res.Plan.Year), Valid: true}
|
||
}
|
||
if res.Plan.Confidence != 0 {
|
||
rec.Confidence = sql.NullFloat64{Float64: res.Plan.Confidence, Valid: true}
|
||
}
|
||
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
log.Warn("recognition reload download failed", "error", err)
|
||
return
|
||
}
|
||
if d.State != store.StateRecognizing {
|
||
// За время вызова LLM задачу увели (cancel/defer) — результат не нужен.
|
||
log.Info("recognition result discarded", "reason", "state_changed", "state", d.State)
|
||
return
|
||
}
|
||
// Claim-токен (MINOR-8): состояние снова recognizing, но за время вызова LLM
|
||
// задачу могли увести из recognizing и вернуть обратно (cancel → relink revive).
|
||
// Тогда это УЖЕ другой recognizing-эпизод (updated_at сменился), а наш результат
|
||
// принадлежит прежней попытке — отбрасываем. Задача остаётся в recognizing, и
|
||
// поллинг-цикл (recognizePending) перезапустит распознавание свежей попыткой.
|
||
//
|
||
// Токен — updated_at секундной точности; сравниваем на строгое равенство, поэтому
|
||
// отбрасываем при ЛЮБОй его смене. В recognizing-эпизоде метку бьют только переходы
|
||
// состояния (не-переходные мутации задачи в этом состоянии не идут), так что валидный
|
||
// результат ложно не теряется; а редкий холостой сброс безопасен — распознавание
|
||
// просто повторится. Остаточное окно — revive целиком в ту же секунду, что и исходный
|
||
// claim (человеческий темп cancel→relink это исключает).
|
||
if d.UpdatedAt != claim {
|
||
log.Info("recognition result discarded", "reason", "stale_claim")
|
||
return
|
||
}
|
||
recID, err := w.store.CreateRecognition(ctx, rec, res.Decision.Reasons)
|
||
if err != nil {
|
||
log.Error("recognition persist failed", "error", err)
|
||
return
|
||
}
|
||
// recognition_id — ключ корреляции попытки (грепается и голым ULID).
|
||
log.Info("recognition persisted", "recognition_id", recID,
|
||
"provider", provider, "provider_id", providerID)
|
||
// Кандидаты базы — для ручного выбора в review.
|
||
if cands := toStoreCandidates(recID, res.Candidates); len(cands) > 0 {
|
||
if err := w.store.CreateCandidates(ctx, cands); err != nil {
|
||
log.Warn("recognition persist candidates failed", "error", err)
|
||
}
|
||
}
|
||
|
||
// Авто-раскладка при подтверждённом матче и чистой валидации (Ф4);
|
||
// иначе — review. Раскладчик может быть не сконфигурирован. При ручной
|
||
// перепривязке (force_review) авто-раскладку не делаем — нужно явное
|
||
// подтверждение человеком.
|
||
overrides := w.overridesOrNil(ctx, id)
|
||
forceReview := overrides[ovrForceReview] == "1"
|
||
if res.Decision.Auto && !forceReview && w.layouter != nil {
|
||
plan := applyOverrides(res.Plan, overrides)
|
||
lctx := w.scoped(ctx, capFileLayout, id, d.PrimaryInfohash())
|
||
// Claim перехода в linking должен закоммититься до хардлинков (MINOR-7):
|
||
// при провале записи не линкуем — задача остаётся в recognizing, и
|
||
// поллинг-цикл (recognizePending) повторит распознавание/авто-раскладку.
|
||
// Правило сходимости папки (и уход в review при рассинхроне) — внутри
|
||
// linkPlan, из состояния linking.
|
||
if err := w.transitionErr(lctx, *d, store.StateLinking, "", ""); err != nil {
|
||
logctx.From(lctx).Warn("auto-apply claim failed, left for recognizing", "error", err)
|
||
return
|
||
}
|
||
if err := w.linkPlan(lctx, d, plan, provider, providerID, savePath); err != nil {
|
||
logctx.From(lctx).Warn("auto-apply failed, left for review", "error", err)
|
||
}
|
||
return
|
||
}
|
||
w.transition(ctx, *d, store.StateReview, "", "")
|
||
}
|
||
|
||
// overridesOrNil читает правки, проглатывая ошибку (для авто-пути).
|
||
func (w *Worker) overridesOrNil(ctx context.Context, id string) map[string]string {
|
||
o, err := w.store.ListOverrides(ctx, id)
|
||
if err != nil {
|
||
logctx.From(ctx).Warn("recognition list overrides failed", "error", err)
|
||
return nil
|
||
}
|
||
return o
|
||
}
|
||
|
||
// --- Команды ревью ---
|
||
|
||
// Apply создаёт хардлинки по текущему плану (с применёнными правками) и
|
||
// переводит задачу в done. Коллизия цели → остаёмся в review с причиной.
|
||
func (w *Worker) Apply(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "apply", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
if w.layouter == nil {
|
||
return fmt.Errorf("apply: layouter not configured")
|
||
}
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("apply: %w", err)
|
||
}
|
||
if d.State != store.StateReview && d.State != store.StateDeferred {
|
||
return fmt.Errorf("apply: download %s is in state %s (expected review/deferred): %w", id, d.State, ErrConflict)
|
||
}
|
||
ctx = w.scoped(ctx, capFileLayout, id, d.PrimaryInfohash())
|
||
|
||
plan, prov, pid, err := w.effectivePlan(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("apply: %w", err)
|
||
}
|
||
t, ok, err := w.torrentByInfohash(ctx, d.HashList())
|
||
if err != nil {
|
||
return fmt.Errorf("apply: lookup torrent: %w", err)
|
||
}
|
||
if !ok {
|
||
// Источник исчез между ревью и применением — приводим состояние к
|
||
// реальности и отказываем (preflight, не доверяем state в БД).
|
||
w.reconcileToReality(ctx, *d, false)
|
||
return fmt.Errorf("apply: источник удалён из qBittorrent: %w", ErrConflict)
|
||
}
|
||
if classify(t.State) != classReady {
|
||
// Последний рубеж перед хардлинками: источник ещё качается — не
|
||
// линкуем неполные файлы (задача могла войти в review иным путём).
|
||
return fmt.Errorf("apply: торрент ещё качается: %w", ErrNotReady)
|
||
}
|
||
|
||
// Claim перехода в linking ОБЯЗАН закоммититься до создания хардлинков: при
|
||
// провале записи не линкуем (иначе ссылки лягут при задаче в review, а
|
||
// финальный linking→done граф отклонит — MINOR-7). Задача остаётся в
|
||
// review/deferred, повтор безопасен. Правило сходимости папки (и уход в
|
||
// review при рассинхроне) — уже внутри linkPlan, из состояния linking.
|
||
if err := w.transitionErr(ctx, *d, store.StateLinking, "", ""); err != nil {
|
||
return fmt.Errorf("apply: %w", err)
|
||
}
|
||
if err := w.linkPlan(ctx, d, plan, prov, pid, translatePath(t.SavePath, w.cfg.PathMap)); err != nil {
|
||
return fmt.Errorf("apply: %w", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// linkPlan строит и создаёт хардлинки по плану, фиксирует батч ссылок и
|
||
// двигает задачу: done при успехе, review при коллизии/невалидном плане/
|
||
// рассинхроне папок, failed при иной ошибке ФС. Идемпотентен (повтор доводит
|
||
// начатое). Под mu; вызывается уже в состоянии linking (claim закоммичен) —
|
||
// поэтому уход в review здесь легален (linking→review), как для коллизии.
|
||
func (w *Worker) linkPlan(ctx context.Context, d *store.Download, plan recognize.Plan, provider, providerID, savePath string) error {
|
||
// Правило сходимости папки: при подтверждённом матче база наследуется от
|
||
// живого якоря; рассинхрон (несколько разных живых папок) → review.
|
||
folderBase, desync, err := w.resolveFolderBase(ctx, d.ID, provider, providerID, layout.MediaType(plan.Type))
|
||
if err != nil {
|
||
w.transition(ctx, *d, store.StateReview, reasonResolve, err.Error())
|
||
return fmt.Errorf("link plan: %w", err)
|
||
}
|
||
if desync {
|
||
w.transition(ctx, *d, store.StateReview, reasonTitleFolderDesync, "несколько живых папок тайтла с одним матчем")
|
||
return fmt.Errorf("рассинхрон папок тайтла: %w", ErrConflict)
|
||
}
|
||
|
||
links, err := w.layouter.BuildLinks(toLayoutPlan(plan, savePath, providerTag(provider, providerID), folderBase))
|
||
if err != nil {
|
||
w.transition(ctx, *d, store.StateReview, reasonBuild, err.Error())
|
||
return fmt.Errorf("build links: %w", err)
|
||
}
|
||
|
||
batch := w.newID()
|
||
results, applyErr := w.layouter.Apply(ctx, links)
|
||
|
||
// Фиксируем то, что успели слинковать (идемпотентность повторного apply).
|
||
fl := make([]store.FileLink, 0, len(results))
|
||
for _, r := range results {
|
||
fl = append(fl, store.FileLink{
|
||
DownloadID: d.ID,
|
||
ApplyBatchID: batch,
|
||
SrcPath: r.Link.Src,
|
||
DstPath: r.Link.Dst,
|
||
Kind: string(r.Link.Kind),
|
||
Status: string(r.Status),
|
||
Size: r.Size,
|
||
})
|
||
}
|
||
if len(fl) > 0 {
|
||
if err := w.store.CreateFileLinks(ctx, fl); err != nil {
|
||
// Хардлинки уже на диске, но их учёт не записан (транзиентная ошибка
|
||
// SQLite). НЕ оставляем задачу в linking (осиротела бы до sweep, а
|
||
// файлы висели бы без file_link — MAJOR-4): уводим в review с
|
||
// причиной. Повторный Apply идемпотентен — Apply вернёт StatusExists
|
||
// на уже созданных ссылках и допишет учёт.
|
||
w.transition(ctx, *d, store.StateReview, reasonPersist, err.Error())
|
||
return fmt.Errorf("persist links: %w", err)
|
||
}
|
||
// Инвариант «один целевой путь — один владелец»: забираем владение
|
||
// фактически разложенными путями у прежних загрузок (см.
|
||
// state-reconciliation). Сверка тех загрузок перестанет считать эти
|
||
// пути своей целью и не «воскресит» их. Это учётная операция, не
|
||
// безопасность данных: файлы уже разложены, поэтому её сбой не
|
||
// стрэндит задачу в linking — логируем WARN и доводим до done, а
|
||
// рассинхрон чужих задач исправит следующий тик сверки.
|
||
owned := make([]string, 0, len(fl))
|
||
for _, l := range fl {
|
||
if isLaidOut(l.Status) {
|
||
owned = append(owned, l.DstPath)
|
||
}
|
||
}
|
||
if err := w.store.SupersedeForeignLinks(ctx, d.ID, owned); err != nil {
|
||
logctx.From(ctx).Warn("supersede foreign links failed", "error", err)
|
||
}
|
||
}
|
||
|
||
if applyErr != nil {
|
||
if errors.Is(applyErr, layout.ErrCollision) {
|
||
w.transition(ctx, *d, store.StateReview, reasonCollision, applyErr.Error())
|
||
return applyErr
|
||
}
|
||
w.transition(ctx, *d, store.StateFailed, "apply", applyErr.Error())
|
||
return applyErr
|
||
}
|
||
|
||
w.transition(ctx, *d, store.StateDone, "", "")
|
||
logctx.From(ctx).Info("layout linked", "batch_id", batch, "links", len(fl))
|
||
return nil
|
||
}
|
||
|
||
// Relink повторно привязывает откатанную (reverted) или отклонённую
|
||
// (cancelled) задачу: возвращает её на распознавание, и поллинг-цикл
|
||
// перезапустит recognize. Авто-раскладку при этом не делаем — ручная
|
||
// перепривязка всегда проходит через ревью с подтверждением (force_review).
|
||
// Источник (раздача в qBittorrent) для этого должен быть на месте и докачан.
|
||
func (w *Worker) Relink(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "relink", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("relink: %w", err)
|
||
}
|
||
if d.State != store.StateReverted && d.State != store.StateCancelled && d.State != store.StateTargetMissing {
|
||
return fmt.Errorf("relink: download %s is in state %s (expected reverted/cancelled/target_missing): %w", id, d.State, ErrConflict)
|
||
}
|
||
// Источник нужен для распознавания и должен быть докачан — проверяем
|
||
// синхронно (без дебаунса); отсутствие приводит состояние к реальности
|
||
// (orphaned/deleted), недокачанный — ErrNotReady.
|
||
if err := w.ensureSourceReady(ctx, d, "relink"); err != nil {
|
||
return err
|
||
}
|
||
// Ручная перепривязка — всегда с подтверждением, без авто-раскладки.
|
||
if err := w.store.SetOverride(ctx, id, ovrForceReview, "1"); err != nil {
|
||
return fmt.Errorf("relink: %w", err)
|
||
}
|
||
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
|
||
// Возврат в активную обработку — только через атомарный гард инварианта
|
||
// «не более одной активной задачи на infohash» (см. design ulid-identity, D4).
|
||
if err := w.store.ActivateIfNoOtherActive(ctx, id, store.StateRecognizing, "", ""); err != nil {
|
||
if errors.Is(err, store.ErrInfohashTaken) {
|
||
return fmt.Errorf("relink: для этого торрента уже есть активная задача: %w", ErrConflict)
|
||
}
|
||
return fmt.Errorf("relink: %w", err)
|
||
}
|
||
logctx.From(ctx).Info("state transition", "from", d.State, "to", store.StateRecognizing)
|
||
return nil
|
||
}
|
||
|
||
// Rerecognize перезапускает распознавание для задачи в review/deferred без
|
||
// добавления подсказки: контекст и прежние подсказки уже накоплены. Поллинг-
|
||
// цикл проведёт задачу recognizing → review заново.
|
||
func (w *Worker) Rerecognize(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "rerecognize", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "rerecognize")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if err := w.ensureSourceReady(ctx, d, "rerecognize"); err != nil {
|
||
return err
|
||
}
|
||
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
|
||
logctx.From(ctx).Info("review re-recognizing without hint")
|
||
w.transition(ctx, *d, store.StateRecognizing, "", "")
|
||
return nil
|
||
}
|
||
|
||
// Refine добавляет подсказку и отправляет задачу на перераспознавание.
|
||
func (w *Worker) Refine(ctx context.Context, id string, hint string) (err error) {
|
||
defer func() { w.logCmd(ctx, "refine", id, err) }()
|
||
hint = strings.TrimSpace(hint)
|
||
if hint == "" {
|
||
return fmt.Errorf("refine: empty hint: %w", ErrInvalidInput)
|
||
}
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "refine")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if err := w.ensureSourceReady(ctx, d, "refine"); err != nil {
|
||
return err
|
||
}
|
||
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
|
||
if err := w.store.AddHint(ctx, id, hint); err != nil {
|
||
return fmt.Errorf("refine: %w", err)
|
||
}
|
||
logctx.From(ctx).Info("review hint added", "hint", hint)
|
||
w.transition(ctx, *d, store.StateRecognizing, "", "")
|
||
return nil
|
||
}
|
||
|
||
// SetType фиксирует тип (override) и перезапускает распознавание с подсказкой
|
||
// — чтобы LLM пересобрал роли файлов под новый тип.
|
||
func (w *Worker) SetType(ctx context.Context, id string, mediaType string) (err error) {
|
||
defer func() { w.logCmd(ctx, "set_type", id, err) }()
|
||
if mediaType != string(recognize.MediaMovie) && mediaType != string(recognize.MediaSeries) {
|
||
return fmt.Errorf("set type: invalid type %q: %w", mediaType, ErrInvalidInput)
|
||
}
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "set type")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if err := w.ensureSourceReady(ctx, d, "set type"); err != nil {
|
||
return err
|
||
}
|
||
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
|
||
if err := w.store.SetOverride(ctx, id, ovrMediaType, mediaType); err != nil {
|
||
return fmt.Errorf("set type: %w", err)
|
||
}
|
||
label := "фильм"
|
||
if mediaType == string(recognize.MediaSeries) {
|
||
label = "сериал"
|
||
}
|
||
if err := w.store.AddHint(ctx, id, "Тип точно: "+label+"."); err != nil {
|
||
return fmt.Errorf("set type: %w", err)
|
||
}
|
||
w.transition(ctx, *d, store.StateRecognizing, "", "")
|
||
return nil
|
||
}
|
||
|
||
// IgnoreFile помечает файл к игнорированию (не линкуем). Остаёмся в review;
|
||
// превью пересчитается с учётом правки.
|
||
func (w *Worker) IgnoreFile(ctx context.Context, id string, src string) (err error) {
|
||
defer func() { w.logCmd(ctx, "ignore_file", id, err) }()
|
||
src = strings.TrimSpace(src)
|
||
if src == "" {
|
||
return fmt.Errorf("ignore: empty path: %w", ErrInvalidInput)
|
||
}
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "ignore")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
overrides, err := w.store.ListOverrides(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("ignore: %w", err)
|
||
}
|
||
ignored := parseIgnored(overrides[ovrIgnoredFiles])
|
||
if !contains(ignored, src) {
|
||
ignored = append(ignored, src)
|
||
}
|
||
b, _ := json.Marshal(ignored)
|
||
if err := w.store.SetOverride(ctx, id, ovrIgnoredFiles, string(b)); err != nil {
|
||
return fmt.Errorf("ignore: %w", err)
|
||
}
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("review file ignored", "src", src)
|
||
return nil
|
||
}
|
||
|
||
// Defer паркует задачу в deferred (вернётся в ревью по действию).
|
||
func (w *Worker) Defer(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "defer", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("defer: %w", err)
|
||
}
|
||
if d.State.IsTerminal() {
|
||
return fmt.Errorf("defer: download %s is terminal (%s): %w", id, d.State, ErrConflict)
|
||
}
|
||
ctx = w.scoped(ctx, capReview, id, d.PrimaryInfohash())
|
||
w.transition(ctx, *d, store.StateDeferred, "", "")
|
||
return nil
|
||
}
|
||
|
||
// Undo снимает хардлинки последнего батча и переводит задачу в reverted.
|
||
// Источник недосягаем (раскладчик удаляет только пути под библиотекой). Откат
|
||
// снимает ЛИШНИЙ хардлинк, а не последнюю копию: layout.Undo отказывается
|
||
// удалять ссылку, если источник уже пропал (nlink<=1) — см. state-reconciliation.
|
||
func (w *Worker) Undo(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "undo", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
if w.layouter == nil {
|
||
return fmt.Errorf("undo: layouter not configured")
|
||
}
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("undo: %w", err)
|
||
}
|
||
if d.State == store.StateOrphaned {
|
||
// Источник удалён → библиотечный хардлинк остался единственной копией.
|
||
// Откат сотрёт данные насовсем — отказываем явно.
|
||
return fmt.Errorf("undo: источник удалён, цель — последняя копия данных, откат невозможен: %w", ErrConflict)
|
||
}
|
||
if d.State != store.StateDone {
|
||
return fmt.Errorf("undo: download %s is in state %s (expected done): %w", id, d.State, ErrConflict)
|
||
}
|
||
ctx = w.scoped(ctx, capFileLayout, id, d.PrimaryInfohash())
|
||
batch, err := w.store.LatestBatchID(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("undo: %w", err)
|
||
}
|
||
if batch == "" {
|
||
return fmt.Errorf("undo: nothing to revert: %w", ErrConflict)
|
||
}
|
||
rows, err := w.store.ListFileLinksByBatch(ctx, batch)
|
||
if err != nil {
|
||
return fmt.Errorf("undo: %w", err)
|
||
}
|
||
n, err := w.layouter.Undo(ctx, laidOutLinks(rows))
|
||
if err != nil {
|
||
return fmt.Errorf("undo: %w", err)
|
||
}
|
||
if err := w.store.DeleteFileLinksByBatch(ctx, batch); err != nil {
|
||
return fmt.Errorf("undo: %w", err)
|
||
}
|
||
w.transition(ctx, *d, store.StateReverted, "", "")
|
||
logctx.From(ctx).Info("layout reverted", "batch_id", batch, "removed", n)
|
||
return nil
|
||
}
|
||
|
||
// laidOutLinks отбирает реально разложенные нами ссылки батча и маппит их в
|
||
// layout.Link. superseded-строки пропускаем: путь забрала другая загрузка (см.
|
||
// state-reconciliation, владение путём), файл по нему — её хардлинк, трогать
|
||
// нельзя. Общий для Undo (снимает с гардом) и Delete (снимает без гарда).
|
||
func laidOutLinks(rows []store.FileLink) []layout.Link {
|
||
links := make([]layout.Link, 0, len(rows))
|
||
for _, r := range rows {
|
||
if !isLaidOut(r.Status) {
|
||
continue
|
||
}
|
||
links = append(links, layout.Link{Src: r.SrcPath, Dst: r.DstPath, Kind: layout.Kind(r.Kind)})
|
||
}
|
||
return links
|
||
}
|
||
|
||
// Delete — пользовательское удаление загрузки: снять наши библиотечные ссылки
|
||
// (гард последней копии ВЫКЛЮЧЕН, в отличие от Undo) И снести раздачу с файлами
|
||
// из qBittorrent, переведя задачу в терминальный deleted. Осознанный выход за
|
||
// инвариант «источник неприкосновенен» — вызывается только после подтверждения
|
||
// в транспорте. Доступно из done/orphaned/target_missing; идемпотентно к
|
||
// отсутствующей стороне. Source-preflight НЕ делает (цель — снять источник,
|
||
// его отсутствие трактуем как уже снятую сторону).
|
||
func (w *Worker) Delete(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "delete", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
if w.layouter == nil {
|
||
return fmt.Errorf("delete: layouter not configured")
|
||
}
|
||
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("delete: %w", err)
|
||
}
|
||
switch d.State {
|
||
case store.StateDone, store.StateOrphaned, store.StateTargetMissing:
|
||
default:
|
||
return fmt.Errorf("delete: download %s is in state %s (expected done/orphaned/target_missing): %w", id, d.State, ErrConflict)
|
||
}
|
||
ctx = w.scoped(ctx, capFileLayout, id, d.PrimaryInfohash())
|
||
|
||
// (а) Снять цель — наши библиотечные ссылки последнего батча БЕЗ гарда
|
||
// последней копии. superseded пропускаем: путь забрала другая загрузка (см.
|
||
// state-reconciliation, владение путём), её хардлинк трогать нельзя. В
|
||
// target_missing / после ручного удаления живых ссылок нет — снятие
|
||
// идемпотентно.
|
||
batch, err := w.store.LatestBatchID(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("delete: %w", err)
|
||
}
|
||
removed := 0
|
||
if batch != "" {
|
||
rows, err := w.store.ListFileLinksByBatch(ctx, batch)
|
||
if err != nil {
|
||
return fmt.Errorf("delete: %w", err)
|
||
}
|
||
removed, err = w.layouter.Remove(ctx, laidOutLinks(rows))
|
||
if err != nil {
|
||
return fmt.Errorf("delete: %w", err)
|
||
}
|
||
if err := w.store.DeleteFileLinksByBatch(ctx, batch); err != nil {
|
||
return fmt.Errorf("delete: %w", err)
|
||
}
|
||
}
|
||
|
||
// (б) Снять источник — раздачу с файлами из qBittorrent. Идемпотентно:
|
||
// отсутствие раздачи (orphaned) не ошибка — qBit не найдёт хеш и ответит OK.
|
||
// Реальную ошибку API пробрасываем и в deleted НЕ уводим: не заявляем
|
||
// освобождённое место, которого не произошло (цель уже снята → сверка
|
||
// приведёт запись к target_missing; повторный delete идемпотентно дожмёт).
|
||
if hashes := d.HashList(); len(hashes) > 0 {
|
||
if err := w.qbt.Delete(ctx, hashes, true); err != nil {
|
||
return fmt.Errorf("delete: qbittorrent: %w", err)
|
||
}
|
||
}
|
||
|
||
// (в) Терминальный deleted с пользовательским маркером инициатора
|
||
// (отличает от reconcile-deleted, который кладёт "reconcile").
|
||
w.transition(ctx, *d, store.StateDeleted, "user_delete", "удалено пользователем")
|
||
// Запись физического эффекта (снятые ссылки) сверх перехода: from/code уже в
|
||
// каноническом `state transition` выше — здесь только отличительное поле.
|
||
logctx.From(ctx).Info("download deleted by user", "removed_links", removed)
|
||
return nil
|
||
}
|
||
|
||
// requireReviewable проверяет, что задача в review/deferred. Вызывается под mu.
|
||
func (w *Worker) requireReviewable(ctx context.Context, id string, op string) (*store.Download, error) {
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("%s: %w", op, err)
|
||
}
|
||
if d.State != store.StateReview && d.State != store.StateDeferred {
|
||
return nil, fmt.Errorf("%s: download %s is in state %s (expected review/deferred): %w", op, id, d.State, ErrConflict)
|
||
}
|
||
return d, nil
|
||
}
|
||
|
||
// --- Выбор базы метаданных (пиннинг; остаёмся в review, применяет человек) ---
|
||
|
||
// ChooseCandidate пиннит выбранного кандидата базы как override (провайдер,
|
||
// id, каноническое имя/год). Раскладку не запускает — превью обновится, а
|
||
// человек подтвердит «Применить».
|
||
func (w *Worker) ChooseCandidate(ctx context.Context, id, candidateID string) (err error) {
|
||
defer func() { w.logCmd(ctx, "choose_candidate", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "choose candidate")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
rec, err := w.store.GetCurrentRecognition(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("choose candidate: %w", err)
|
||
}
|
||
cand, err := w.store.GetCandidate(ctx, candidateID)
|
||
if err != nil {
|
||
return fmt.Errorf("choose candidate: %w", err)
|
||
}
|
||
if rec == nil || cand == nil || cand.RecognitionID != rec.ID {
|
||
return fmt.Errorf("choose candidate: candidate %s does not belong to the current recognition: %w", candidateID, ErrInvalidInput)
|
||
}
|
||
return w.chooseCandidateLocked(ctx, id, d, rec, *cand)
|
||
}
|
||
|
||
// AddManualSource добавляет источник вручную по (provider, id) и выбирает его.
|
||
// Когда автопоиск промахнулся: сохраняем кандидата (дедуп по provider:id) и
|
||
// пиннит как выбранный. provider — из набора tmdb/tvdb/imdb.
|
||
func (w *Worker) AddManualSource(ctx context.Context, id, provider, providerID string) (err error) {
|
||
defer func() { w.logCmd(ctx, "add_manual_source", id, err) }()
|
||
provider = strings.TrimSpace(strings.ToLower(provider))
|
||
providerID = strings.TrimSpace(providerID)
|
||
switch provider {
|
||
case "tmdb", "tvdb", "imdb":
|
||
default:
|
||
return fmt.Errorf("add source: invalid provider %q (tmdb/tvdb/imdb): %w", provider, ErrInvalidInput)
|
||
}
|
||
if providerID == "" {
|
||
return fmt.Errorf("add source: empty id: %w", ErrInvalidInput)
|
||
}
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "add source")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
rec, err := w.store.GetCurrentRecognition(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("add source: %w", err)
|
||
}
|
||
if rec == nil {
|
||
return fmt.Errorf("add source: no recognition")
|
||
}
|
||
cand, err := w.findOrCreateCandidate(ctx, rec.ID, provider, providerID)
|
||
if err != nil {
|
||
return fmt.Errorf("add source: %w", err)
|
||
}
|
||
return w.chooseCandidateLocked(ctx, id, d, rec, *cand)
|
||
}
|
||
|
||
// findOrCreateCandidate возвращает кандидата рекогниции по (provider, id),
|
||
// создавая его при отсутствии (дедуп по provider:id). Под mu.
|
||
func (w *Worker) findOrCreateCandidate(ctx context.Context, recognitionID, provider, providerID string) (*store.MetadataCandidate, error) {
|
||
cands, err := w.store.ListCandidatesByRecognition(ctx, recognitionID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if c := findCandidate(cands, provider, providerID); c != nil {
|
||
return c, nil
|
||
}
|
||
if err := w.store.CreateCandidates(ctx, []store.MetadataCandidate{{
|
||
RecognitionID: recognitionID,
|
||
Provider: provider,
|
||
ProviderID: providerID,
|
||
}}); err != nil {
|
||
return nil, err
|
||
}
|
||
cands, err = w.store.ListCandidatesByRecognition(ctx, recognitionID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if c := findCandidate(cands, provider, providerID); c != nil {
|
||
return c, nil
|
||
}
|
||
return nil, fmt.Errorf("candidate %s:%s not found after create", provider, providerID)
|
||
}
|
||
|
||
func findCandidate(cands []store.MetadataCandidate, provider, providerID string) *store.MetadataCandidate {
|
||
for i := range cands {
|
||
if cands[i].Provider == provider && cands[i].ProviderID == providerID {
|
||
return &cands[i]
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// chooseCandidateLocked пиннит кандидата как выбранный источник. Пишет ПОЛНЫЙ
|
||
// самосогласованный набор пинов (provider/id/title/year): пустые title/year у
|
||
// кандидата очищают возможный унаследованный пин прежнего источника — иначе
|
||
// превью разошлось бы с применением (решение 1a). Под mu.
|
||
func (w *Worker) chooseCandidateLocked(ctx context.Context, id string, d *store.Download, rec *store.Recognition, cand store.MetadataCandidate) error {
|
||
title := ""
|
||
if cand.Title.Valid {
|
||
title = cand.Title.String
|
||
}
|
||
year := 0
|
||
if cand.Year.Valid {
|
||
year = int(cand.Year.Int64)
|
||
}
|
||
// Режиссёр выбранного источника из метабазы (credits) — best-effort: пустой
|
||
// (нет провайдера/не умеет/ошибка) очищает возможный унаследованный пин; тогда
|
||
// в ярлыке режиссёр берётся из распознавания/авто-матча (plan.Director), а если
|
||
// и там пусто — из сохранённого контекста (нижний слой).
|
||
director := ""
|
||
if w.recognizer != nil {
|
||
director = w.recognizer.Director(ctx, candMediaType(rec), cand.Provider, cand.ProviderID)
|
||
}
|
||
for field, value := range sourcePins(cand.Provider, cand.ProviderID, title, year, director) {
|
||
if err := w.store.SetOverride(ctx, id, field, value); err != nil {
|
||
return fmt.Errorf("choose candidate: %w", err)
|
||
}
|
||
}
|
||
if err := w.store.SetCandidateChosen(ctx, rec.ID, cand.ID); err != nil {
|
||
return fmt.Errorf("choose candidate: %w", err)
|
||
}
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("review candidate chosen",
|
||
"provider", cand.Provider, "provider_id", cand.ProviderID)
|
||
// Подтверждённый матч — переливаем каноническое имя в display_name и в ярлык
|
||
// раздачи (best-effort, косметика). Сбой обновления имени не должен ронять
|
||
// выбор кандидата: логируем и продолжаем.
|
||
if err := w.refreshDisplayNameLocked(ctx, id); err != nil {
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).
|
||
Warn("display name refresh after candidate choice failed", "error", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// SetProviderID пиннит провайдера и id вручную (без выбора из списка).
|
||
func (w *Worker) SetProviderID(ctx context.Context, id string, provider, providerID string) (err error) {
|
||
defer func() { w.logCmd(ctx, "set_provider_id", id, err) }()
|
||
provider = strings.TrimSpace(strings.ToLower(provider))
|
||
providerID = strings.TrimSpace(providerID)
|
||
switch provider {
|
||
case "tmdb", "tvdb", "imdb":
|
||
default:
|
||
return fmt.Errorf("set provider: invalid provider %q (tmdb/tvdb/imdb): %w", provider, ErrInvalidInput)
|
||
}
|
||
if providerID == "" {
|
||
return fmt.Errorf("set provider: empty id: %w", ErrInvalidInput)
|
||
}
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "set provider")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// Режиссёр вручную заданного источника из метабазы (credits) — best-effort
|
||
// косметика: сбой чтения рекогниции не валит смену источника (тип по умолчанию
|
||
// movie, как трактует candMediaType(nil)).
|
||
director := ""
|
||
if w.recognizer != nil {
|
||
rec, rerr := w.store.GetCurrentRecognition(ctx, id)
|
||
if rerr != nil {
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).
|
||
Warn("set provider: recognition lookup for director failed", "error", rerr)
|
||
rec = nil
|
||
}
|
||
director = w.recognizer.Director(ctx, candMediaType(rec), provider, providerID)
|
||
}
|
||
// Полный набор пинов: id задан вручную, название/год берём из плана
|
||
// (очищаем возможный унаследованный пин прежнего источника).
|
||
for field, value := range sourcePins(provider, providerID, "", 0, director) {
|
||
if err := w.store.SetOverride(ctx, id, field, value); err != nil {
|
||
return fmt.Errorf("set provider: %w", err)
|
||
}
|
||
}
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("review provider set",
|
||
"provider", provider, "provider_id", providerID)
|
||
return nil
|
||
}
|
||
|
||
// ClearProvider — «без базы»: снимает матч (тег папки не ставится) и очищает
|
||
// пины названия/года (источник — распознавание нейронкой).
|
||
func (w *Worker) ClearProvider(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "clear_provider", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
d, err := w.requireReviewable(ctx, id, "clear provider")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
for field, value := range sourcePins("none", "", "", 0, "") {
|
||
if err := w.store.SetOverride(ctx, id, field, value); err != nil {
|
||
return fmt.Errorf("clear provider: %w", err)
|
||
}
|
||
}
|
||
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("review provider cleared")
|
||
return nil
|
||
}
|
||
|
||
// sourcePins — полный самосогласованный набор пинов источника (решение 1a):
|
||
// title/year пишутся пустой строкой, если у источника их нет; в applyOverrides
|
||
// пустая строка трактуется как «нет override» → берётся значение плана. Так
|
||
// выбор любого источника даёт детерминированный эффективный план, а превью
|
||
// совпадает с применением. Используется и в коммите (SetOverride), и в превью.
|
||
func sourcePins(provider, providerID, title string, year int, director string) map[string]string {
|
||
yr := ""
|
||
if year > 0 {
|
||
yr = strconv.Itoa(year)
|
||
}
|
||
return map[string]string{
|
||
ovrProvider: provider,
|
||
ovrProviderID: providerID,
|
||
ovrTitle: title,
|
||
ovrYear: yr,
|
||
ovrDirector: director,
|
||
}
|
||
}
|
||
|
||
// candMediaType — тип из текущей рекогниции (для выборки режиссёра нужным
|
||
// эндпоинтом провайдера). Неизвестный/пустой → movie (дефолт схемы).
|
||
func candMediaType(rec *store.Recognition) recognize.MediaType {
|
||
if rec != nil && rec.MediaType.String == string(recognize.MediaSeries) {
|
||
return recognize.MediaSeries
|
||
}
|
||
return recognize.MediaMovie
|
||
}
|
||
|
||
// --- Данные для экрана ревью ---
|
||
|
||
// ReviewData — всё, что нужно транспорту для отрисовки ревью.
|
||
type ReviewData struct {
|
||
Download store.Download
|
||
Recognition *store.Recognition
|
||
Plan recognize.Plan // эффективный (с применёнными правками)
|
||
Preview []layout.Link // целевые пути активного источника (Src — относительный)
|
||
Candidates []store.MetadataCandidate // кандидаты базы для ручного выбора
|
||
Sources []SourceOption // единый список источников совпадения (нейронка + кандидаты)
|
||
Provider string // эффективный провайдер (с учётом выбора)
|
||
ProviderID string // эффективный id в базе
|
||
Hints []string
|
||
Overrides map[string]string
|
||
}
|
||
|
||
// SourceKind — вид источника в едином списке ревью.
|
||
type SourceKind string
|
||
|
||
const (
|
||
SourceNeural SourceKind = "neural" // распознавание нейронкой (без базы)
|
||
SourceCandidate SourceKind = "candidate" // кандидат метабазы (в т.ч. добавленный вручную)
|
||
)
|
||
|
||
// SourceOption — источник совпадения в списке ревью: эффективные поля и
|
||
// предпросмотр целевых путей, посчитанные эфемерно (без записи overrides).
|
||
type SourceOption struct {
|
||
Kind SourceKind
|
||
CandidateID string // ULID кандидата (пусто для нейронки)
|
||
Provider string // "none" для нейронки
|
||
ProviderID string
|
||
URL string // ссылка кандидата на запись (если есть)
|
||
Title string // эффективное название для этого источника
|
||
Year int // эффективный год
|
||
Type string // "movie" | "series"
|
||
Active bool // текущий эффективный источник
|
||
Plan recognize.Plan // эффективный план (для показа файлов → раскладка)
|
||
Preview []layout.Link // целевые пути этого источника
|
||
}
|
||
|
||
// ReviewData собирает данные ревью по загрузке.
|
||
func (w *Worker) ReviewData(ctx context.Context, id string) (*ReviewData, error) {
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("review data: %w", err)
|
||
}
|
||
log := logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash()))
|
||
rec, err := w.store.GetCurrentRecognition(ctx, id)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("review data: %w", err)
|
||
}
|
||
hints, err := w.store.ListHints(ctx, id)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("review data: %w", err)
|
||
}
|
||
overrides, err := w.store.ListOverrides(ctx, id)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("review data: %w", err)
|
||
}
|
||
|
||
prov, pid := effectiveProvider(rec, overrides)
|
||
rd := &ReviewData{
|
||
Download: *d, Recognition: rec, Hints: hints, Overrides: overrides,
|
||
Provider: prov, ProviderID: pid,
|
||
}
|
||
if rec != nil {
|
||
if cands, cerr := w.store.ListCandidatesByRecognition(ctx, rec.ID); cerr == nil {
|
||
rd.Candidates = cands
|
||
} else {
|
||
log.Debug("review data list candidates failed", "error", cerr)
|
||
}
|
||
}
|
||
if rec != nil && rec.Plan.Valid {
|
||
var rawPlan recognize.Plan
|
||
if err := json.Unmarshal([]byte(rec.Plan.String), &rawPlan); err != nil {
|
||
log.Warn("review data unmarshal plan failed", "error", err)
|
||
} else {
|
||
rd.Plan = applyOverrides(rawPlan, overrides)
|
||
// Превью активного источника строим по относительным путям с
|
||
// provider-тегом; ошибку логируем на Debug — покажем без превью.
|
||
// База наследуется тем же правилом сходимости, что и применение
|
||
// (превью=применение); рассинхрон в превью не переводит в review.
|
||
if w.layouter != nil {
|
||
tag := providerTag(prov, pid)
|
||
base, _, berr := w.resolveFolderBase(ctx, id, prov, pid, layout.MediaType(rd.Plan.Type))
|
||
if berr != nil {
|
||
log.Debug("review data resolve folder base failed", "error", berr)
|
||
}
|
||
if links, lerr := w.layouter.BuildLinks(toLayoutPlan(rd.Plan, "", tag, base)); lerr == nil {
|
||
rd.Preview = links
|
||
} else {
|
||
// Видимая деградация: без превью на экране ревью пропадает
|
||
// кнопка «Применить» — не рядовой Debug, а WARN.
|
||
log.Warn("review data build preview failed", "error", lerr)
|
||
}
|
||
}
|
||
// Единый список источников: нейронка + кандидаты, каждый с
|
||
// эфемерным превью из сырого плана (без записи overrides).
|
||
rd.Sources = w.buildSources(ctx, id, rawPlan, overrides, prov, pid, rd.Candidates)
|
||
}
|
||
}
|
||
return rd, nil
|
||
}
|
||
|
||
// buildSources собирает единый список источников: нейронка (первой) +
|
||
// кандидаты (дедуп по provider:id). Активным помечается текущий эффективный
|
||
// источник.
|
||
func (w *Worker) buildSources(ctx context.Context, downloadID string, rawPlan recognize.Plan, overrides map[string]string, prov, pid string, cands []store.MetadataCandidate) []SourceOption {
|
||
neutral := prov == "" || prov == "none"
|
||
out := make([]SourceOption, 0, len(cands)+1)
|
||
out = append(out, w.sourceOption(ctx, downloadID, SourceNeural, rawPlan, overrides, "", "none", "", "", "", 0, neutral))
|
||
seen := map[string]bool{}
|
||
for _, c := range cands {
|
||
key := c.Provider + ":" + c.ProviderID
|
||
if seen[key] {
|
||
continue
|
||
}
|
||
seen[key] = true
|
||
title := ""
|
||
if c.Title.Valid {
|
||
title = c.Title.String
|
||
}
|
||
year := 0
|
||
if c.Year.Valid {
|
||
year = int(c.Year.Int64)
|
||
}
|
||
active := !neutral && c.Provider == prov && c.ProviderID == pid
|
||
out = append(out, w.sourceOption(ctx, downloadID, SourceCandidate, rawPlan, overrides, c.ID, c.Provider, c.ProviderID, c.URL.String, title, year, active))
|
||
}
|
||
return out
|
||
}
|
||
|
||
// sourceOption строит один источник: накладывает его пины на неисточниковые
|
||
// overrides, считает эффективный план и предпросмотр путей — эфемерно, без
|
||
// записи. Гарантия preview == apply: тот же набор пинов запишет выбор.
|
||
func (w *Worker) sourceOption(ctx context.Context, downloadID string, kind SourceKind, rawPlan recognize.Plan, base map[string]string, candID, provider, providerID, url, title string, year int, active bool) SourceOption {
|
||
// Режиссёр в эфемерном превью не тянем (пусто): он не влияет на пути раскладки,
|
||
// а per-candidate выборка credits дорога. Режиссёр появляется в инфо-панели
|
||
// после закрепления выбора (ovrDirector персистится в chooseCandidateLocked).
|
||
eff := applyOverrides(rawPlan, mergeSourceOverrides(base, sourcePins(provider, providerID, title, year, "")))
|
||
opt := SourceOption{
|
||
Kind: kind,
|
||
CandidateID: candID,
|
||
Provider: provider,
|
||
ProviderID: providerID,
|
||
URL: url,
|
||
Title: eff.Title,
|
||
Year: eff.Year,
|
||
Type: string(eff.Type),
|
||
Active: active,
|
||
Plan: eff,
|
||
}
|
||
if w.layouter != nil {
|
||
// Превью источника наследует базу тем же правилом сходимости, что и
|
||
// применение (превью=применение); рассинхрон в превью не переводит в review.
|
||
folderBase, _, _ := w.resolveFolderBase(ctx, downloadID, provider, providerID, layout.MediaType(eff.Type))
|
||
if links, err := w.layouter.BuildLinks(toLayoutPlan(eff, "", providerTag(provider, providerID), folderBase)); err == nil {
|
||
opt.Preview = links
|
||
}
|
||
}
|
||
return opt
|
||
}
|
||
|
||
// mergeSourceOverrides накладывает пины источника (provider/id/title/year/
|
||
// director) на неисточниковые overrides (media_type, ignored_files, force_review,
|
||
// ...).
|
||
func mergeSourceOverrides(base, pins map[string]string) map[string]string {
|
||
m := make(map[string]string, len(base)+len(pins))
|
||
for k, v := range base {
|
||
switch k {
|
||
case ovrProvider, ovrProviderID, ovrTitle, ovrYear, ovrDirector:
|
||
continue
|
||
default:
|
||
m[k] = v
|
||
}
|
||
}
|
||
maps.Copy(m, pins)
|
||
return m
|
||
}
|
||
|
||
// effectivePlan загружает текущий план, применяет правки и возвращает
|
||
// эффективные provider/provider_id (для тега папки и правила сходимости) (под mu).
|
||
func (w *Worker) effectivePlan(ctx context.Context, id string) (plan recognize.Plan, provider, providerID string, err error) {
|
||
rec, err := w.store.GetCurrentRecognition(ctx, id)
|
||
if err != nil {
|
||
return recognize.Plan{}, "", "", err
|
||
}
|
||
if rec == nil || !rec.Plan.Valid {
|
||
return recognize.Plan{}, "", "", fmt.Errorf("no recognition plan")
|
||
}
|
||
if err := json.Unmarshal([]byte(rec.Plan.String), &plan); err != nil {
|
||
return recognize.Plan{}, "", "", fmt.Errorf("parse plan: %w", err)
|
||
}
|
||
overrides, err := w.store.ListOverrides(ctx, id)
|
||
if err != nil {
|
||
return recognize.Plan{}, "", "", err
|
||
}
|
||
prov, pid := effectiveProvider(rec, overrides)
|
||
return applyOverrides(plan, overrides), prov, pid, nil
|
||
}
|
||
|
||
// effectiveDisplayName собирает полный ярлык display_name из эффективных полей.
|
||
// Тонкая обёртка над naming.EffectiveFields (слоистое разрешение план →
|
||
// контекст) + Label — единый источник разрешения полей ярлыка, общий со
|
||
// страницей загрузки и экраном ревью.
|
||
func effectiveDisplayName(d store.Download, plan recognize.Plan) string {
|
||
return naming.EffectiveFields(d.ParsedContext, plan).Label()
|
||
}
|
||
|
||
// RefreshDisplayName — внешняя точка входа обновления отображаемого имени
|
||
// (ручная кнопка): берёт w.mu и делегирует refreshDisplayNameLocked.
|
||
func (w *Worker) RefreshDisplayName(ctx context.Context, id string) (err error) {
|
||
defer func() { w.logCmd(ctx, "refresh_name", id, err) }()
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
return w.refreshDisplayNameLocked(ctx, id)
|
||
}
|
||
|
||
// refreshDisplayNameLocked переливает уже вычисленное каноническое имя
|
||
// (эффективный план с учётом пинов) в download.display_name и в имя раздачи
|
||
// qBittorrent — без нового вызова LLM. Косметика: не влияет на пути/раскладку.
|
||
// Best-effort к qBittorrent: недоступность/отсутствие раздачи не проваливает
|
||
// операцию (display_name пишется в любом случае). Вызывается под w.mu.
|
||
func (w *Worker) refreshDisplayNameLocked(ctx context.Context, id string) error {
|
||
d, err := w.store.GetDownload(ctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("refresh display name: %w", err)
|
||
}
|
||
if d == nil {
|
||
return fmt.Errorf("refresh display name: download %s: %w", id, store.ErrNotFound)
|
||
}
|
||
cctx := w.scoped(ctx, capIngest, id, d.PrimaryInfohash())
|
||
|
||
rec, err := w.store.GetCurrentRecognition(cctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("refresh display name: %w", err)
|
||
}
|
||
if rec == nil || !rec.Plan.Valid {
|
||
return nil // распознавания ещё нет — обновлять нечего (no-op)
|
||
}
|
||
plan, _, _, err := w.effectivePlan(cctx, id)
|
||
if err != nil {
|
||
return fmt.Errorf("refresh display name: %w", err)
|
||
}
|
||
name := effectiveDisplayName(*d, plan)
|
||
if name == "" {
|
||
return nil // пустое эффективное название — no-op
|
||
}
|
||
if err := w.store.SetDisplayName(cctx, id, name); err != nil {
|
||
return fmt.Errorf("refresh display name: %w", err)
|
||
}
|
||
logctx.From(cctx).Info("display name refreshed", "display_name", name)
|
||
|
||
// Переименование раздачи — best-effort по её реальному Torrent.Hash (для
|
||
// гибрид/v2 он может не совпасть с нашим primary infohash). Отсутствие
|
||
// раздачи (удалена) — штатный no-op: display_name уже обновлён.
|
||
t, ok, err := w.torrentByInfohash(cctx, d.HashList())
|
||
if err != nil {
|
||
logctx.From(cctx).Warn("refresh display name: qbittorrent lookup failed", "error", err)
|
||
return nil
|
||
}
|
||
if !ok {
|
||
return nil
|
||
}
|
||
if err := w.qbt.RenameTorrent(cctx, t.Hash, name); err != nil {
|
||
logctx.From(cctx).Warn("refresh display name: qbittorrent rename failed", "error", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Коды причины (error_code) ухода задачи в review при раскладке (linkPlan) —
|
||
// корреляционный ключ шага, на котором раскладка остановилась. Свод в одном
|
||
// месте (как errCode* в worker.go); человекочитаемый текст кладётся в error_msg.
|
||
const (
|
||
reasonResolve = "resolve" // не удалось разрешить базу папки тайтла
|
||
reasonBuild = "build" // не удалось построить план ссылок
|
||
reasonPersist = "persist" // ссылки на диске, но учёт не записан
|
||
reasonCollision = "collision" // целевой путь уже занят (layout.ErrCollision)
|
||
reasonTitleFolderDesync = "title_folder_desync" // ≥2 разных живых папок тайтла с одним матчем
|
||
)
|
||
|
||
// resolveFolderBase применяет правило сходимости папки (см. file-layout spec):
|
||
// при подтверждённом матче наследует базу имени от живой папки-якоря того же
|
||
// (provider, providerID), кроме самой загрузки downloadID. Возвращает базу для
|
||
// layout.Plan.FolderBase (пусто → печатать из распознавания), флаг рассинхрона
|
||
// (≥2 разных живых папок) и ошибку. Чистая (только чтение БД+ФС), без побочных
|
||
// эффектов — переиспользуется применением и превью. Живость якоря определяется
|
||
// существованием папки на диске (os.Lstat), а не статусом ссылки в БД.
|
||
func (w *Worker) resolveFolderBase(ctx context.Context, downloadID, provider, providerID string, mt layout.MediaType) (base string, desync bool, err error) {
|
||
if provider == "" || provider == "none" || providerID == "" || w.layouter == nil {
|
||
return "", false, nil
|
||
}
|
||
paths, err := w.store.LiveTitleFolders(ctx, provider, providerID, downloadID)
|
||
if err != nil {
|
||
return "", false, fmt.Errorf("resolve folder base: %w", err)
|
||
}
|
||
bases := make(map[string]string, 2) // абсолютная папка тайтла → база имени
|
||
for _, p := range paths {
|
||
dir, b, ok := w.layouter.TitleFolder(mt, p)
|
||
if !ok {
|
||
continue // не под корнем / не разобрать — не якорь
|
||
}
|
||
if _, seen := bases[dir]; seen {
|
||
continue
|
||
}
|
||
if _, serr := os.Lstat(dir); serr != nil {
|
||
continue // папки на диске нет (переименована/удалена) — не якорь
|
||
}
|
||
bases[dir] = b
|
||
}
|
||
switch len(bases) {
|
||
case 0:
|
||
return "", false, nil
|
||
case 1:
|
||
for _, b := range bases {
|
||
return b, false, nil
|
||
}
|
||
}
|
||
return "", true, nil // рассинхрон: несколько разных живых папок
|
||
}
|
||
|
||
// --- Хелперы преобразования ---
|
||
|
||
// applyOverrides применяет ручные правки к плану: форсит тип, каноническое
|
||
// имя/год (из выбранного кандидата базы) и помечает игнорируемые файлы ролью
|
||
// ignore (их раскладка пропустит).
|
||
func applyOverrides(plan recognize.Plan, overrides map[string]string) recognize.Plan {
|
||
if mt := overrides[ovrMediaType]; mt == string(recognize.MediaMovie) || mt == string(recognize.MediaSeries) {
|
||
plan.Type = recognize.MediaType(mt)
|
||
}
|
||
if t := overrides[ovrTitle]; t != "" {
|
||
plan.Title = t
|
||
}
|
||
if y := overrides[ovrYear]; y != "" {
|
||
if year, err := strconv.Atoi(y); err == nil {
|
||
plan.Year = year
|
||
}
|
||
}
|
||
if d := overrides[ovrDirector]; d != "" {
|
||
plan.Director = d
|
||
}
|
||
ignored := parseIgnored(overrides[ovrIgnoredFiles])
|
||
if len(ignored) > 0 {
|
||
for i := range plan.Files {
|
||
if contains(ignored, plan.Files[i].Src) {
|
||
plan.Files[i].Role = "ignore"
|
||
}
|
||
}
|
||
}
|
||
return plan
|
||
}
|
||
|
||
// effectiveProvider возвращает провайдера и id для тега папки с учётом
|
||
// ручного выбора: запиненный override перекрывает распознанный матч.
|
||
// override "none" означает явный отказ от базы.
|
||
func effectiveProvider(rec *store.Recognition, overrides map[string]string) (provider, id string) {
|
||
if p, ok := overrides[ovrProvider]; ok {
|
||
return p, overrides[ovrProviderID]
|
||
}
|
||
if rec != nil {
|
||
return rec.Provider.String, rec.ProviderID.String
|
||
}
|
||
return "", ""
|
||
}
|
||
|
||
// toStoreCandidates переводит кандидатов распознавания в строки БД,
|
||
// подставляя тег-предпочтительный provider/id (внешний из TVMaze и т.п.).
|
||
func toStoreCandidates(recognitionID string, cands []metadata.Candidate) []store.MetadataCandidate {
|
||
out := make([]store.MetadataCandidate, 0, len(cands))
|
||
for _, c := range cands {
|
||
prov, id := recognize.CandidateTag(c)
|
||
mc := store.MetadataCandidate{
|
||
RecognitionID: recognitionID,
|
||
Provider: prov,
|
||
ProviderID: id,
|
||
Title: store.NullString(c.Title),
|
||
}
|
||
if c.Year != 0 {
|
||
mc.Year = sql.NullInt64{Int64: int64(c.Year), Valid: true}
|
||
}
|
||
if c.URL != "" {
|
||
mc.URL = store.NullString(c.URL)
|
||
}
|
||
out = append(out, mc)
|
||
}
|
||
return out
|
||
}
|
||
|
||
// ProviderTag — экспорт providerTag для диагностических команд (CLI
|
||
// `jellybit recognize --dry-run`).
|
||
func ProviderTag(provider, id string) string { return providerTag(provider, id) }
|
||
|
||
// ToLayoutPlan — экспорт toLayoutPlan для диагностических команд.
|
||
func ToLayoutPlan(p recognize.Plan, srcPrefix, providerTag string) layout.Plan {
|
||
return toLayoutPlan(p, srcPrefix, providerTag, "")
|
||
}
|
||
|
||
// providerTag строит тег папки для Jellyfin из провайдера и id: "tmdbid-…"
|
||
// / "tvdbid-…". Пустой id (нет матча) → пустой тег.
|
||
func providerTag(provider, id string) string {
|
||
if id == "" {
|
||
return ""
|
||
}
|
||
switch provider {
|
||
case "tmdb":
|
||
return "tmdbid-" + id
|
||
case "tvdb":
|
||
return "tvdbid-" + id
|
||
case "imdb":
|
||
return "imdbid-" + id
|
||
default:
|
||
return ""
|
||
}
|
||
}
|
||
|
||
// toLayoutPlan переводит план распознавания в план раскладки. srcPrefix
|
||
// (savePath) приклеивается к относительным путям файлов; пустой — оставляет
|
||
// относительные (для превью). providerTag добавляется к имени папки. folderBase
|
||
// (непустой) — унаследованная от живого якоря база имени (правило сходимости):
|
||
// перекрывает Title/Year в папке и в именах файлов. Роли вне
|
||
// main/episode/subtitle отбрасываются.
|
||
func toLayoutPlan(plan recognize.Plan, srcPrefix, providerTag, folderBase string) layout.Plan {
|
||
lp := layout.Plan{
|
||
Type: layout.MediaType(plan.Type),
|
||
Title: plan.Title,
|
||
Year: plan.Year,
|
||
ProviderTag: providerTag,
|
||
FolderBase: folderBase,
|
||
}
|
||
for _, f := range plan.Files {
|
||
role, ok := mapRole(f.Role)
|
||
if !ok {
|
||
continue
|
||
}
|
||
src := f.Src
|
||
if srcPrefix != "" {
|
||
src = filepath.Join(srcPrefix, f.Src)
|
||
}
|
||
lp.Files = append(lp.Files, layout.PlanFile{
|
||
Src: src,
|
||
Role: role,
|
||
Season: f.Season,
|
||
Episode: f.Episode,
|
||
})
|
||
}
|
||
return lp
|
||
}
|
||
|
||
func mapRole(r recognize.FileRole) (layout.Role, bool) {
|
||
switch r {
|
||
case recognize.RoleMain:
|
||
return layout.RoleMain, true
|
||
case recognize.RoleEpisode:
|
||
return layout.RoleEpisode, true
|
||
case recognize.RoleSubtitle:
|
||
return layout.RoleSubtitle, true
|
||
default:
|
||
return "", false
|
||
}
|
||
}
|
||
|
||
// torrentByInfohash ищет торрент по любому из хешей загрузки (v1/v2/hash).
|
||
// Листаем ВСЕ торренты (а не только свою категорию): раздача могла быть
|
||
// усыновлена по тегу и иметь чужую/пустую категорию — фильтр по категории её
|
||
// бы потерял (как и в Poll, см. там же).
|
||
func (w *Worker) torrentByInfohash(ctx context.Context, hashes []string) (qbt.Torrent, bool, error) {
|
||
torrents, err := w.qbt.Torrents(ctx, "")
|
||
if err != nil {
|
||
return qbt.Torrent{}, false, err
|
||
}
|
||
want := make(map[string]bool, len(hashes))
|
||
for _, h := range hashes {
|
||
want[store.NormalizeHash(h)] = true
|
||
}
|
||
for _, t := range torrents {
|
||
for _, h := range torrentIndexHashes(t) {
|
||
if want[strings.ToLower(h)] {
|
||
return t, true, nil
|
||
}
|
||
}
|
||
}
|
||
return qbt.Torrent{}, false, nil
|
||
}
|
||
|
||
func parseIgnored(s string) []string {
|
||
if s == "" {
|
||
return nil
|
||
}
|
||
var out []string
|
||
// Ошибку разбора глотаем намеренно: битый JSON в override ignored_files
|
||
// (не должен возникать — пишем его сами через json.Marshal) трактуем как
|
||
// «нет игнора», а не роняем команду. Худший исход — файл не будет пропущен,
|
||
// человек увидит его в превью и пометит заново.
|
||
_ = json.Unmarshal([]byte(s), &out)
|
||
return out
|
||
}
|
||
|
||
func contains(ss []string, s string) bool {
|
||
for _, x := range ss {
|
||
if x == s {
|
||
return true
|
||
}
|
||
}
|
||
return false
|
||
}
|