package worker import ( "context" "database/sql" "encoding/json" "errors" "fmt" "maps" "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/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" // запиненный год 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, "", "") } 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, 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 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, tag := "none", "", "" if res.Match != nil { provider, providerID = res.Match.Provider, res.Match.ProviderID tag = providerTag(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", "state", d.State) 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()) w.transition(lctx, *d, store.StateLinking, "", "") if err := w.linkPlan(lctx, d, plan, tag, 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) error { 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, tag, 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) } w.transition(ctx, *d, store.StateLinking, "", "") if err := w.linkPlan(ctx, d, plan, tag, translatePath(t.SavePath, w.cfg.PathMap)); err != nil { return fmt.Errorf("apply: %w", err) } return nil } // linkPlan строит и создаёт хардлинки по плану, фиксирует батч ссылок и // двигает задачу: done при успехе, review при коллизии/невалидном плане, // failed при иной ошибке ФС. Идемпотентен (повтор доводит начатое). Под mu. func (w *Worker) linkPlan(ctx context.Context, d *store.Download, plan recognize.Plan, providerTag, savePath string) error { links, err := w.layouter.BuildLinks(toLayoutPlan(plan, savePath, providerTag)) if err != nil { w.transition(ctx, *d, store.StateReview, "build", 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 { 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, "collision", 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) error { 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). if err := w.ensureSourcePresent(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("relink re-recognizing", "from", d.State) return nil } // Rerecognize перезапускает распознавание для задачи в review/deferred без // добавления подсказки: контекст и прежние подсказки уже накоплены. Поллинг- // цикл проведёт задачу recognizing → review заново. func (w *Worker) Rerecognize(ctx context.Context, id string) error { w.mu.Lock() defer w.mu.Unlock() d, err := w.requireReviewable(ctx, id, "rerecognize") if err != nil { return err } if err := w.ensureSourcePresent(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) error { hint = strings.TrimSpace(hint) if hint == "" { return fmt.Errorf("refine: empty hint") } w.mu.Lock() defer w.mu.Unlock() d, err := w.requireReviewable(ctx, id, "refine") if err != nil { return err } if err := w.ensureSourcePresent(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) error { if mediaType != string(recognize.MediaMovie) && mediaType != string(recognize.MediaSeries) { return fmt.Errorf("set type: invalid type %q", mediaType) } w.mu.Lock() defer w.mu.Unlock() d, err := w.requireReviewable(ctx, id, "set type") if err != nil { return err } if err := w.ensureSourcePresent(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) error { src = strings.TrimSpace(src) if src == "" { return fmt.Errorf("ignore: empty path") } 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) error { 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)", id, d.State) } 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) error { 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") } rows, err := w.store.ListFileLinksByBatch(ctx, batch) if err != nil { return fmt.Errorf("undo: %w", err) } // Снимаем только реально разложенные нами ссылки. superseded-строки — // путь забрала другая загрузка (см. state-reconciliation, владение // путём); файл по нему теперь её хардлинк, трогать его нельзя. 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)}) } n, err := w.layouter.Undo(ctx, links) 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 } // 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) error { 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", candidateID) } 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) error { 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)", provider) } if providerID == "" { return fmt.Errorf("add source: empty id") } 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) } for field, value := range sourcePins(cand.Provider, cand.ProviderID, title, year) { 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) return nil } // SetProviderID пиннит провайдера и id вручную (без выбора из списка). func (w *Worker) SetProviderID(ctx context.Context, id string, provider, providerID string) error { 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)", provider) } if providerID == "" { return fmt.Errorf("set provider: empty id") } w.mu.Lock() defer w.mu.Unlock() d, err := w.requireReviewable(ctx, id, "set provider") if err != nil { return err } // Полный набор пинов: id задан вручную, название/год берём из плана // (очищаем возможный унаследованный пин прежнего источника). for field, value := range sourcePins(provider, providerID, "", 0) { 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) error { 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) map[string]string { yr := "" if year > 0 { yr = strconv.Itoa(year) } return map[string]string{ ovrProvider: provider, ovrProviderID: providerID, ovrTitle: title, ovrYear: yr, } } // --- Данные для экрана ревью --- // 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 — покажем без превью. if w.layouter != nil { tag := providerTag(prov, pid) if links, lerr := w.layouter.BuildLinks(toLayoutPlan(rd.Plan, "", tag)); lerr == nil { rd.Preview = links } else { log.Debug("review data build preview failed", "error", lerr) } } // Единый список источников: нейронка + кандидаты, каждый с // эфемерным превью из сырого плана (без записи overrides). rd.Sources = w.buildSources(rawPlan, overrides, prov, pid, rd.Candidates) } } return rd, nil } // buildSources собирает единый список источников: нейронка (первой) + // кандидаты (дедуп по provider:id). Активным помечается текущий эффективный // источник. func (w *Worker) buildSources(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(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(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(kind SourceKind, rawPlan recognize.Plan, base map[string]string, candID, provider, providerID, url, title string, year int, active bool) SourceOption { 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 { if links, err := w.layouter.BuildLinks(toLayoutPlan(eff, "", providerTag(provider, providerID))); err == nil { opt.Preview = links } } return opt } // mergeSourceOverrides накладывает пины источника (provider/id/title/year) на // неисточниковые 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: continue default: m[k] = v } } maps.Copy(m, pins) return m } // effectivePlan загружает текущий план, применяет правки и возвращает // provider-тег для имени папки (под mu). func (w *Worker) effectivePlan(ctx context.Context, id string) (recognize.Plan, string, 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") } var plan recognize.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), providerTag(prov, pid), 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 } } 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 добавляется к имени папки. Роли // вне main/episode/subtitle отбрасываются. func toLayoutPlan(plan recognize.Plan, srcPrefix, providerTag string) layout.Plan { lp := layout.Plan{ Type: layout.MediaType(plan.Type), Title: plan.Title, Year: plan.Year, ProviderTag: providerTag, } 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 []string{t.Hash, t.InfohashV1, t.InfohashV2} { if h != "" && 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.Unmarshal([]byte(s), &out) return out } func contains(ss []string, s string) bool { for _, x := range ss { if x == s { return true } } return false }