Обработка рассинхрона состояния с реальностью (state-reconciliation)
Распознаём ручное удаление источника (раздача в qBittorrent) и/или цели (разложенные хардлинки) и отражаем его в состоянии задачи, без автодействий. - Новая capability state-reconciliation (OpenSpec): фоновая сверка по матрице «источник × цель» → состояния target_missing/orphaned/deleted, переходы и самовосстановление (healing). - worker: reconcileDesync в Poll (только разложенные/desync-задачи), дебаунс пропажи источника (порог [worker].source_missing_threshold) и синхронный preflight перед действиями (relink/recognize/apply/undo) — не доверяем state в БД. - layout.Undo: отказ снять последнюю копию (nlink<=1 или нет источника), отказ всего батча без частичного отката (ErrLastCopy). - store: единый список terminalStates для IsTerminal и FindActiveByInfohash (иначе семантика «активности» разъезжается), столбец source_miss_count, миграция 0003. - httpapi/web и Telegram: показ новых состояний и уведомления о рассинхроне. - Доки: workflow.md, jellyfin-layout.md, database.md (+0003), config. Change заархивирован в openspec/changes/archive, дельта влита в openspec/specs. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,196 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/layout"
|
||||
"git.vakhrushev.me/av/jellybit/internal/logctx"
|
||||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
)
|
||||
|
||||
// desyncStates — состояния, которые ведёт сверка с реальностью: уже
|
||||
// разложенные (done) и сами состояния рассинхрона. Активные и
|
||||
// пользовательски-терминальные (reverted/cancelled/failed/stuck) сверка не
|
||||
// трогает (см. state-reconciliation).
|
||||
var desyncStates = []store.State{
|
||||
store.StateDone,
|
||||
store.StateTargetMissing,
|
||||
store.StateOrphaned,
|
||||
store.StateDeleted,
|
||||
}
|
||||
|
||||
// deriveState выводит состояние задачи из двумерной матрицы «источник × цель»
|
||||
// (см. state-reconciliation, D1).
|
||||
func deriveState(sourcePresent, targetPresent bool) store.State {
|
||||
switch {
|
||||
case sourcePresent && targetPresent:
|
||||
return store.StateDone
|
||||
case sourcePresent && !targetPresent:
|
||||
return store.StateTargetMissing
|
||||
case !sourcePresent && targetPresent:
|
||||
return store.StateOrphaned
|
||||
default:
|
||||
return store.StateDeleted
|
||||
}
|
||||
}
|
||||
|
||||
// reconcileDesync сверяет разложенные/desync-задачи с реальностью. Вызывается
|
||||
// из Poll под w.mu. byHash — карта присутствующих в qBittorrent раздач.
|
||||
func (w *Worker) reconcileDesync(ctx context.Context, byHash map[string]qbt.Torrent) {
|
||||
cands, err := w.store.ListDownloadsByState(ctx, desyncStates...)
|
||||
if err != nil {
|
||||
w.log.Warn("reconcile list desync failed", "capability", capIngest, "error", err)
|
||||
return
|
||||
}
|
||||
for _, d := range cands {
|
||||
w.reconcileOneDesync(ctx, d, byHash)
|
||||
}
|
||||
}
|
||||
|
||||
// reconcileOneDesync сверяет одну задачу: вычисляет присутствие источника (с
|
||||
// дебаунсом) и цели, выводит состояние и переходит при изменении.
|
||||
func (w *Worker) reconcileOneDesync(ctx context.Context, d store.Download, byHash map[string]qbt.Torrent) {
|
||||
if !d.Infohash.Valid {
|
||||
return // нечем сопоставить источник
|
||||
}
|
||||
ctx = w.scoped(ctx, capIngest, d.ID, d.Infohash.String)
|
||||
_, sourceSeen := byHash[strings.ToLower(d.Infohash.String)]
|
||||
|
||||
// Дебаунс пропажи источника: считаем удалённым только после порога подряд
|
||||
// идущих промахов; любое появление сбрасывает счётчик.
|
||||
sourcePresent := w.debounceSource(ctx, d, sourceSeen)
|
||||
|
||||
targetPresent, err := w.targetPresent(ctx, d.ID)
|
||||
if err != nil {
|
||||
logctx.From(ctx).Warn("reconcile target probe failed", "error", err)
|
||||
return // не можем проверить цель — не дёргаем состояние
|
||||
}
|
||||
|
||||
want := deriveState(sourcePresent, targetPresent)
|
||||
if want == d.State {
|
||||
return
|
||||
}
|
||||
w.transition(ctx, d, want, "reconcile", reconcileReason(sourcePresent, targetPresent))
|
||||
}
|
||||
|
||||
// debounceSource обновляет счётчик промахов источника и возвращает, считать ли
|
||||
// источник присутствующим с учётом порога. Запись счётчика — только при его
|
||||
// изменении (без лишних UPDATE на каждом тике).
|
||||
func (w *Worker) debounceSource(ctx context.Context, d store.Download, sourceSeen bool) bool {
|
||||
threshold := max(w.cfg.SourceMissingThreshold, 1)
|
||||
if sourceSeen {
|
||||
if d.SourceMissCount != 0 {
|
||||
if err := w.store.SetSourceMissCount(ctx, d.ID, 0); err != nil {
|
||||
logctx.From(ctx).Warn("reconcile reset miss count failed", "error", err)
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
miss := d.SourceMissCount + 1
|
||||
if err := w.store.SetSourceMissCount(ctx, d.ID, miss); err != nil {
|
||||
logctx.From(ctx).Warn("reconcile bump miss count failed", "error", err)
|
||||
}
|
||||
// До порога источник трактуется как присутствующий — задача не дёргается.
|
||||
return miss < threshold
|
||||
}
|
||||
|
||||
// targetPresent сообщает, существуют ли разложенные хардлинки задачи. Цель
|
||||
// считается присутствующей, только если существуют ВСЕ ссылки последнего
|
||||
// батча; частичная пропажа — это отсутствие цели (библиотека сломана → relink).
|
||||
func (w *Worker) targetPresent(ctx context.Context, id int64) (bool, error) {
|
||||
batch, err := w.store.LatestBatchID(ctx, id)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("latest batch: %w", err)
|
||||
}
|
||||
if batch == "" {
|
||||
return false, nil // ничего не разложено — цели нет
|
||||
}
|
||||
rows, err := w.store.ListFileLinksByBatch(ctx, batch)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("list links: %w", err)
|
||||
}
|
||||
laidOut := 0
|
||||
for _, r := range rows {
|
||||
// Все статусы, означающие реальный файл на диске: только что слинкован,
|
||||
// скопирован (фолбэк) или уже существовал тем же inode (идемпотентный
|
||||
// повтор apply). Прочие (collision и т.п.) — не наша разложенная цель.
|
||||
if !isLaidOut(r.Status) {
|
||||
continue
|
||||
}
|
||||
laidOut++
|
||||
if _, err := os.Lstat(r.DstPath); err != nil {
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return false, nil
|
||||
}
|
||||
return false, fmt.Errorf("stat target %q: %w", r.DstPath, err)
|
||||
}
|
||||
}
|
||||
return laidOut > 0, nil
|
||||
}
|
||||
|
||||
// isLaidOut сообщает, означает ли статус file_link реально лежащий на ФС файл
|
||||
// нашей раскладки (а не коллизию или иной не-разложенный исход).
|
||||
func isLaidOut(status string) bool {
|
||||
switch layout.LinkStatus(status) {
|
||||
case layout.StatusLinked, layout.StatusCopied, layout.StatusExists:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// reconcileReason — человекочитаемая причина перехода для лога/UI.
|
||||
func reconcileReason(sourcePresent, targetPresent bool) string {
|
||||
switch {
|
||||
case sourcePresent && targetPresent:
|
||||
return "источник и цель на месте"
|
||||
case sourcePresent && !targetPresent:
|
||||
return "цель удалена, источник на месте"
|
||||
case !sourcePresent && targetPresent:
|
||||
return "источник удалён, цель (последняя копия) на месте"
|
||||
default:
|
||||
return "удалены и источник, и цель"
|
||||
}
|
||||
}
|
||||
|
||||
// --- Синхронный preflight перед действием (не доверяем state в БД) ---
|
||||
|
||||
// ensureSourcePresent синхронно (без дебаунса) проверяет, что раздача есть в
|
||||
// qBittorrent прямо сейчас. При отсутствии приводит состояние к реальности и
|
||||
// возвращает ErrConflict. Недоступность qBittorrent — честный отказ операции.
|
||||
func (w *Worker) ensureSourcePresent(ctx context.Context, d *store.Download, op string) error {
|
||||
if !d.Infohash.Valid {
|
||||
return fmt.Errorf("%s: download %d has no infohash", op, d.ID)
|
||||
}
|
||||
_, ok, err := w.torrentByInfohash(ctx, d.Infohash.String)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%s: %w", op, err)
|
||||
}
|
||||
if ok {
|
||||
return nil
|
||||
}
|
||||
// Источник пропал — немедленно приводим состояние к реальности.
|
||||
w.reconcileToReality(ctx, *d, false)
|
||||
return fmt.Errorf("%s: источник удалён из qBittorrent: %w", op, ErrConflict)
|
||||
}
|
||||
|
||||
// reconcileToReality выводит и проставляет состояние по уже известному факту об
|
||||
// источнике (sourcePresent) и фактически проверенной цели. Используется
|
||||
// preflight-проверками: немедленно, без дебаунса.
|
||||
func (w *Worker) reconcileToReality(ctx context.Context, d store.Download, sourcePresent bool) {
|
||||
targetPresent, err := w.targetPresent(ctx, d.ID)
|
||||
if err != nil {
|
||||
logctx.FromOr(ctx, w.log).Warn("preflight target probe failed", "error", err)
|
||||
return
|
||||
}
|
||||
want := deriveState(sourcePresent, targetPresent)
|
||||
if want == d.State {
|
||||
return
|
||||
}
|
||||
w.transition(ctx, d, want, "reconcile", reconcileReason(sourcePresent, targetPresent))
|
||||
}
|
||||
@@ -0,0 +1,186 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
)
|
||||
|
||||
// reconcileFixture готовит done-задачу с одной разложенной ссылкой: цель —
|
||||
// файл в temp-каталоге (создаётся при makeTarget), источник — наличие
|
||||
// раздачи в fakeQbt (sourcePresent).
|
||||
type reconcileFixture struct {
|
||||
w *Worker
|
||||
st *memStore
|
||||
dst string
|
||||
}
|
||||
|
||||
func newReconcileFixture(t *testing.T, state store.State, sourcePresent, makeTarget bool) reconcileFixture {
|
||||
t.Helper()
|
||||
dir := t.TempDir()
|
||||
dst := filepath.Join(dir, "Movie (2024).mkv")
|
||||
src := filepath.Join(dir, "src.mkv")
|
||||
if makeTarget {
|
||||
if err := os.WriteFile(dst, []byte("video"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
st := newMemStore()
|
||||
d := completedDownload(1)
|
||||
d.State = state
|
||||
st.put(d)
|
||||
st.links = append(st.links, store.FileLink{
|
||||
DownloadID: 1, ApplyBatchID: "b1", SrcPath: src, DstPath: dst,
|
||||
Kind: "video", Status: "linked",
|
||||
})
|
||||
|
||||
var torrents []qbt.Torrent
|
||||
if sourcePresent {
|
||||
torrents = []qbt.Torrent{{Hash: ihTest}}
|
||||
}
|
||||
w := testWorkerWith(st, &fakeQbt{torrents: torrents}, nil, nil)
|
||||
w.cfg.SourceMissingThreshold = 1 // помечаем при первой же пропаже (без задержки)
|
||||
return reconcileFixture{w: w, st: st, dst: dst}
|
||||
}
|
||||
|
||||
func TestReconcileMatrix(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
source, target bool
|
||||
want store.State
|
||||
}{
|
||||
{"источник+цель → done", true, true, store.StateDone},
|
||||
{"цель удалена → target_missing", true, false, store.StateTargetMissing},
|
||||
{"источник удалён → orphaned", false, true, store.StateOrphaned},
|
||||
{"оба удалены → deleted", false, false, store.StateDeleted},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
f := newReconcileFixture(t, store.StateDone, tc.source, tc.target)
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != tc.want {
|
||||
t.Errorf("state = %q, want %q", got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestReconcileHealing(t *testing.T) {
|
||||
// Источник вернулся, цель на месте → orphaned лечится обратно в done.
|
||||
f := newReconcileFixture(t, store.StateOrphaned, true, true)
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateDone {
|
||||
t.Errorf("state = %q, want done (healing)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReconcilePartialTargetLoss(t *testing.T) {
|
||||
// Две ссылки, одна цель удалена → цель считается отсутствующей.
|
||||
f := newReconcileFixture(t, store.StateDone, true, true)
|
||||
missing := filepath.Join(filepath.Dir(f.dst), "Movie (2024).en.srt")
|
||||
f.st.links = append(f.st.links, store.FileLink{
|
||||
DownloadID: 1, ApplyBatchID: "b1", SrcPath: "/x.srt", DstPath: missing,
|
||||
Kind: "subtitle", Status: "linked",
|
||||
})
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateTargetMissing {
|
||||
t.Errorf("state = %q, want target_missing (частичная пропажа)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReconcileDebounce(t *testing.T) {
|
||||
// Порог 3: первые два промаха не помечают, третий — помечает orphaned;
|
||||
// возврат источника лечит обратно и сбрасывает счётчик.
|
||||
f := newReconcileFixture(t, store.StateDone, false, true)
|
||||
f.w.cfg.SourceMissingThreshold = 3
|
||||
|
||||
for i := 1; i <= 2; i++ {
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll %d: %v", i, err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateDone {
|
||||
t.Fatalf("tick %d: state = %q, want done (до порога)", i, got)
|
||||
}
|
||||
if got := f.st.downloads[1].SourceMissCount; got != i {
|
||||
t.Errorf("tick %d: miss = %d, want %d", i, got, i)
|
||||
}
|
||||
}
|
||||
if err := f.w.Poll(context.Background()); err != nil { // третий промах
|
||||
t.Fatalf("Poll 3: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateOrphaned {
|
||||
t.Fatalf("tick 3: state = %q, want orphaned (порог достигнут)", got)
|
||||
}
|
||||
|
||||
// Источник вернулся.
|
||||
f.w.qbt.(*fakeQbt).torrents = []qbt.Torrent{{Hash: ihTest}}
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll heal: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateDone {
|
||||
t.Errorf("state = %q, want done (источник вернулся)", got)
|
||||
}
|
||||
if got := f.st.downloads[1].SourceMissCount; got != 0 {
|
||||
t.Errorf("miss = %d, want 0 (сброс)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReconcileSkipsActiveStates(t *testing.T) {
|
||||
// downloading сверкой не трогаем, даже если раздачи нет в qBittorrent.
|
||||
f := newReconcileFixture(t, store.StateDownloading, false, true)
|
||||
if err := f.w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateDownloading {
|
||||
t.Errorf("state = %q, want downloading (сверка не трогает активные)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUndoRejectedForOrphaned(t *testing.T) {
|
||||
f := newReconcileFixture(t, store.StateOrphaned, false, true)
|
||||
err := f.w.Undo(context.Background(), 1)
|
||||
if err == nil {
|
||||
t.Fatal("ожидали отказ Undo для orphaned")
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateOrphaned {
|
||||
t.Errorf("state = %q, want orphaned (без изменений)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRelinkFromTargetMissing(t *testing.T) {
|
||||
// target_missing + источник на месте → relink ведёт в recognizing.
|
||||
f := newReconcileFixture(t, store.StateTargetMissing, true, false)
|
||||
if err := f.w.Relink(context.Background(), 1); err != nil {
|
||||
t.Fatalf("Relink: %v", err)
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateRecognizing {
|
||||
t.Errorf("state = %q, want recognizing", got)
|
||||
}
|
||||
if f.st.overrides[1][ovrForceReview] != "1" {
|
||||
t.Errorf("force_review = %q, want 1", f.st.overrides[1][ovrForceReview])
|
||||
}
|
||||
}
|
||||
|
||||
func TestPreflightFixesStaleState(t *testing.T) {
|
||||
// В БД target_missing (источник якобы есть), но фактически источник пропал,
|
||||
// а цель на месте: relink немедленно приводит состояние к orphaned, не
|
||||
// дожидаясь фоновой сверки.
|
||||
f := newReconcileFixture(t, store.StateTargetMissing, false, true)
|
||||
if err := f.w.Relink(context.Background(), 1); err == nil {
|
||||
t.Fatal("ожидали отказ relink при пропавшем источнике")
|
||||
}
|
||||
if got := f.st.downloads[1].State; got != store.StateOrphaned {
|
||||
t.Errorf("state = %q, want orphaned (preflight привёл к реальности)", got)
|
||||
}
|
||||
}
|
||||
+27
-12
@@ -238,7 +238,10 @@ func (w *Worker) Apply(ctx context.Context, id int64) error {
|
||||
return fmt.Errorf("apply: lookup torrent: %w", err)
|
||||
}
|
||||
if !ok {
|
||||
return fmt.Errorf("apply: torrent not found")
|
||||
// Источник исчез между ревью и применением — приводим состояние к
|
||||
// реальности и отказываем (preflight, не доверяем state в БД).
|
||||
w.reconcileToReality(ctx, *d, false)
|
||||
return fmt.Errorf("apply: источник удалён из qBittorrent: %w", ErrConflict)
|
||||
}
|
||||
|
||||
w.transition(ctx, *d, store.StateLinking, "", "")
|
||||
@@ -306,17 +309,13 @@ func (w *Worker) Relink(ctx context.Context, id int64) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("relink: %w", err)
|
||||
}
|
||||
if d.State != store.StateReverted && d.State != store.StateCancelled {
|
||||
return fmt.Errorf("relink: download %d is in state %s (expected reverted/cancelled): %w", id, d.State, ErrConflict)
|
||||
if d.State != store.StateReverted && d.State != store.StateCancelled && d.State != store.StateTargetMissing {
|
||||
return fmt.Errorf("relink: download %d is in state %s (expected reverted/cancelled/target_missing): %w", id, d.State, ErrConflict)
|
||||
}
|
||||
if !d.Infohash.Valid {
|
||||
return fmt.Errorf("relink: download %d has no infohash", id)
|
||||
}
|
||||
// Раздача должна ещё быть в qBittorrent — без неё распознавать нечего.
|
||||
if _, ok, terr := w.torrentByInfohash(ctx, d.Infohash.String); terr != nil {
|
||||
return fmt.Errorf("relink: %w", terr)
|
||||
} else if !ok {
|
||||
return fmt.Errorf("relink: торрент не найден в qBittorrent")
|
||||
// Источник нужен для распознавания — проверяем синхронно (без дебаунса) и при
|
||||
// его отсутствии приводим состояние к реальности (orphaned/deleted).
|
||||
if err := w.ensureSourcePresent(ctx, d, "relink"); err != nil {
|
||||
return err
|
||||
}
|
||||
// Вернуть задачу в активную обработку можно, только если другой активной
|
||||
// задачи на этот infohash нет (partial unique index по idempotency_key).
|
||||
@@ -348,6 +347,9 @@ func (w *Worker) Rerecognize(ctx context.Context, id int64) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := w.ensureSourcePresent(ctx, d, "rerecognize"); err != nil {
|
||||
return err
|
||||
}
|
||||
ctx = w.scoped(ctx, capReview, id, d.Infohash.String)
|
||||
logctx.From(ctx).Info("review re-recognizing without hint")
|
||||
w.transition(ctx, *d, store.StateRecognizing, "", "")
|
||||
@@ -367,6 +369,9 @@ func (w *Worker) Refine(ctx context.Context, id int64, hint string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := w.ensureSourcePresent(ctx, d, "refine"); err != nil {
|
||||
return err
|
||||
}
|
||||
ctx = w.scoped(ctx, capReview, id, d.Infohash.String)
|
||||
if err := w.store.AddHint(ctx, id, hint); err != nil {
|
||||
return fmt.Errorf("refine: %w", err)
|
||||
@@ -389,6 +394,9 @@ func (w *Worker) SetType(ctx context.Context, id int64, mediaType string) error
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := w.ensureSourcePresent(ctx, d, "set type"); err != nil {
|
||||
return err
|
||||
}
|
||||
ctx = w.scoped(ctx, capReview, id, d.Infohash.String)
|
||||
if err := w.store.SetOverride(ctx, id, ovrMediaType, mediaType); err != nil {
|
||||
return fmt.Errorf("set type: %w", err)
|
||||
@@ -452,7 +460,9 @@ func (w *Worker) Defer(ctx context.Context, id int64) error {
|
||||
}
|
||||
|
||||
// Undo снимает хардлинки последнего батча и переводит задачу в reverted.
|
||||
// Источник недосягаем (раскладчик удаляет только пути под библиотекой).
|
||||
// Источник недосягаем (раскладчик удаляет только пути под библиотекой). Откат
|
||||
// снимает ЛИШНИЙ хардлинк, а не последнюю копию: layout.Undo отказывается
|
||||
// удалять ссылку, если источник уже пропал (nlink<=1) — см. state-reconciliation.
|
||||
func (w *Worker) Undo(ctx context.Context, id int64) error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
@@ -464,6 +474,11 @@ func (w *Worker) Undo(ctx context.Context, id int64) error {
|
||||
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 %d is in state %s (expected done): %w", id, d.State, ErrConflict)
|
||||
}
|
||||
|
||||
@@ -155,7 +155,8 @@ func TestRerecognize_ReviewToRecognizing(t *testing.T) {
|
||||
d := completedDownload(1)
|
||||
d.State = store.StateReview
|
||||
st.put(d)
|
||||
w := testWorkerWith(st, &fakeQbt{}, &fakeRecognizer{}, nil)
|
||||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ihTest}}}
|
||||
w := testWorkerWith(st, qb, &fakeRecognizer{}, nil)
|
||||
|
||||
if err := w.Rerecognize(context.Background(), 1); err != nil {
|
||||
t.Fatalf("Rerecognize: %v", err)
|
||||
@@ -184,8 +185,10 @@ func TestRelink_TorrentMissing(t *testing.T) {
|
||||
if err := w.Relink(context.Background(), 1); err == nil {
|
||||
t.Fatal("ожидали ошибку при отсутствии торрента, получили nil")
|
||||
}
|
||||
if st.downloads[1].State != store.StateReverted {
|
||||
t.Errorf("state = %q, want reverted (без изменений)", st.downloads[1].State)
|
||||
// Preflight приводит состояние к реальности: источника нет и цели нет
|
||||
// (reverted — ссылки сняты) → deleted (см. state-reconciliation).
|
||||
if st.downloads[1].State != store.StateDeleted {
|
||||
t.Errorf("state = %q, want deleted (preflight привёл к реальности)", st.downloads[1].State)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -291,6 +294,13 @@ func (m *memStore) SetDownloadState(_ context.Context, id int64, st store.State,
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *memStore) SetSourceMissCount(_ context.Context, id int64, n int) error {
|
||||
if d, ok := m.downloads[id]; ok {
|
||||
d.SourceMissCount = n
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *memStore) CreateRecognition(_ context.Context, r *store.Recognition, reasons []string) (int64, error) {
|
||||
for _, e := range m.recs {
|
||||
if e.DownloadID == r.DownloadID {
|
||||
@@ -578,7 +588,8 @@ func TestRefine_AddsHintAndRerecognizes(t *testing.T) {
|
||||
d := completedDownload(1)
|
||||
d.State = store.StateReview
|
||||
st.put(d)
|
||||
w := testWorkerWith(st, &fakeQbt{}, &fakeRecognizer{}, nil)
|
||||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ihTest}}}
|
||||
w := testWorkerWith(st, qb, &fakeRecognizer{}, nil)
|
||||
|
||||
if err := w.Refine(context.Background(), 1, "это второй сезон"); err != nil {
|
||||
t.Fatalf("Refine: %v", err)
|
||||
@@ -599,7 +610,8 @@ func TestSetType(t *testing.T) {
|
||||
d := completedDownload(1)
|
||||
d.State = store.StateReview
|
||||
st.put(d)
|
||||
w := testWorkerWith(st, &fakeQbt{}, &fakeRecognizer{}, nil)
|
||||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ihTest}}}
|
||||
w := testWorkerWith(st, qb, &fakeRecognizer{}, nil)
|
||||
|
||||
if err := w.SetType(context.Background(), 1, "series"); err != nil {
|
||||
t.Fatalf("SetType: %v", err)
|
||||
|
||||
@@ -40,6 +40,7 @@ type Store interface {
|
||||
ListDownloadsByState(ctx context.Context, states ...store.State) ([]store.Download, error)
|
||||
GetDownload(ctx context.Context, id int64) (*store.Download, error)
|
||||
SetDownloadState(ctx context.Context, id int64, state store.State, errCode, errMsg string) error
|
||||
SetSourceMissCount(ctx context.Context, id int64, n int) error
|
||||
|
||||
// Discovery (усыновление раздач по категории/тегу).
|
||||
ExistsByInfohash(ctx context.Context, infohash string) (bool, error)
|
||||
@@ -88,8 +89,10 @@ type Layouter interface {
|
||||
type NotifyEvent string
|
||||
|
||||
const (
|
||||
EventReview NotifyEvent = "review" // задача ждёт подтверждения
|
||||
EventDone NotifyEvent = "done" // раскладка завершена
|
||||
EventReview NotifyEvent = "review" // задача ждёт подтверждения
|
||||
EventDone NotifyEvent = "done" // раскладка завершена
|
||||
EventOrphaned NotifyEvent = "orphaned" // источник пропал, цель — последняя копия
|
||||
EventTargetMissing NotifyEvent = "target_missing" // цель удалена, доступен relink
|
||||
)
|
||||
|
||||
// Notifier — исходящие пинги (Telegram). Вызывается неблокирующе.
|
||||
@@ -113,6 +116,9 @@ type Config struct {
|
||||
PollInterval time.Duration
|
||||
StuckAfter time.Duration // stalledDL дольше → stuck
|
||||
MagnetTimeout time.Duration // metaDL дольше → failed
|
||||
// SourceMissingThreshold — порог дебаунса пропажи источника (тиков сверки).
|
||||
// <1 трактуется как 1 (помечаем при первой же устойчивой пропаже).
|
||||
SourceMissingThreshold int
|
||||
}
|
||||
|
||||
// Worker — поллер и владелец переходов.
|
||||
@@ -235,6 +241,10 @@ func (w *Worker) Poll(ctx context.Context) error {
|
||||
}
|
||||
w.reconcile(ctx, d, t)
|
||||
}
|
||||
|
||||
// Сверка разложенных задач с реальностью (источник в qBit + хардлинки на ФС)
|
||||
// — отдельно от активных, по двумерной матрице (см. state-reconciliation).
|
||||
w.reconcileDesync(ctx, byHash)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -293,6 +303,10 @@ func (w *Worker) transition(ctx context.Context, d store.Download, state store.S
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventReview)
|
||||
case store.StateDone:
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventDone)
|
||||
case store.StateOrphaned:
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventOrphaned)
|
||||
case store.StateTargetMissing:
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventTargetMissing)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -90,6 +90,15 @@ func (f *fakeStore) SetDownloadState(_ context.Context, id int64, st store.State
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeStore) SetSourceMissCount(_ context.Context, id int64, n int) error {
|
||||
d, ok := f.downloads[id]
|
||||
if !ok {
|
||||
return fmt.Errorf("download %d not found", id)
|
||||
}
|
||||
d.SourceMissCount = n
|
||||
return nil
|
||||
}
|
||||
|
||||
// --- Ф3-методы Store (заглушки; переопределяются в review_test.go) ---
|
||||
|
||||
func (f *fakeStore) CreateRecognition(_ context.Context, _ *store.Recognition, _ []string) (int64, error) {
|
||||
|
||||
Reference in New Issue
Block a user