Отмена задачи (catched→cancelled) в окно, пока worker вне блокировки выводит имя (LLM) и делает qbt.Add, оставляла добавленный торрент в qBittorrent без задачи-владельца: PromoteCatched корректно пропускал переход, но источник уже качался/сидировал вечно, а усыновить его назад нельзя (хеши принадлежат отменённой задаче). Спека покрывала переход состояния, но не побочный эффект. Комбинированная защита в processCatched: - re-read состояния под w.mu прямо перед qbt.Add — при отмене источник не добавляется вовсе (сужает окно гонки); - свежий листинг перед add подтверждает отсутствие infohash — признак «своего» торрента; при сбое листинга/присутствии add не делаем (усыновит следующий тик); - при отмене в окне после add (промах PromoteCatched, подтверждённый re-read'ом state != catched) — уборка добавленного нами торрента qbt.Delete(_, true); - WARN/ERROR-логи по этому пути с корреляцией по download_id, без секретов. Гарантия «удаляем только своё»: удаление-с-данными достижимо ТОЛЬКО после подтверждённого отсутствия infohash перед add, поэтому пред-существующий/чужой торрент с тем же хешем никогда не сносится (негативный инвариант). Обоснование по инварианту «источник неприкосновенен» — в design.md изменения. Дельта — download-tracking (требование «Добавление пойманной загрузки в qBittorrent»): re-read перед add, подтверждение отсутствия, уборка при отмене, негативный сценарий. Тесты покрывают все ветки (skip-before-add, cleanup после add, пред-существующий не удаляется, сбой БД не удаляет, сбой листинга не добавляет). Change archived: 2026-07-17-cancel-during-add-cleanup. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
363 lines
16 KiB
Go
363 lines
16 KiB
Go
package worker
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"testing"
|
||
"time"
|
||
|
||
"git.vakhrushev.me/av/jellybit/internal/naming"
|
||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||
)
|
||
|
||
// fakeNamer — вывод имени для шага добавления. onCall позволяет вклиниться в
|
||
// момент (медленного) вывода имени, симулируя параллельную отмену.
|
||
type fakeNamer struct {
|
||
name string
|
||
fields naming.Fields // извлечённая структура (ok=true, если Title непуст)
|
||
gotContext string
|
||
calls int
|
||
onCall func()
|
||
}
|
||
|
||
func (f *fakeNamer) Derive(_ context.Context, contextText, _ string) (string, naming.Fields, bool) {
|
||
f.calls++
|
||
f.gotContext = contextText
|
||
if f.onCall != nil {
|
||
f.onCall()
|
||
}
|
||
return f.name, f.fields, f.fields.Title != ""
|
||
}
|
||
|
||
func catchedStore(id, infohash, createdAt, ctxText string) *fakeStore {
|
||
return &fakeStore{downloads: map[string]*store.Download{
|
||
id: {
|
||
ID: id,
|
||
State: store.StateCatched,
|
||
SourceType: store.SourceMagnet,
|
||
SourceRef: "magnet:?xt=urn:btih:" + infohash + "&dn=Dune",
|
||
Infohashes: hashesOf(id, infohash),
|
||
Context: ctxText,
|
||
CreatedAt: createdAt,
|
||
},
|
||
}}
|
||
}
|
||
|
||
const catchedIH = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
|
||
// now воркера в тестах — 2026-06-14 10:00 UTC (см. newTestWorker).
|
||
var nowStr = store.FormatTime(time.Date(2026, 6, 14, 10, 0, 0, 0, time.UTC))
|
||
|
||
// Успех: выводим имя, добавляем в qBit с rename, переводим catched → downloading
|
||
// и сохраняем display_name.
|
||
func TestProcessCatchedAddsToQbit(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "Дюна 2")
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
nm := &fakeNamer{
|
||
name: "Дюна: Часть вторая (2024)",
|
||
fields: naming.Fields{Type: "movie", Title: "Дюна: Часть вторая", Director: "Дени Вильнёв", Year: 2024},
|
||
}
|
||
w.SetNamer(nm)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("qbt.Add calls = %d, want 1", len(qb.added))
|
||
}
|
||
add := qb.added[0]
|
||
if add.Rename != "Дюна: Часть вторая (2024)" {
|
||
t.Errorf("rename = %q", add.Rename)
|
||
}
|
||
if add.Category != "jellybit" || add.URLs[0] != st.downloads["1"].SourceRef {
|
||
t.Errorf("add = %+v", add)
|
||
}
|
||
if nm.gotContext != "Дюна 2" {
|
||
t.Errorf("namer получил контекст %q", nm.gotContext)
|
||
}
|
||
d := st.downloads["1"]
|
||
if d.State != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading", d.State)
|
||
}
|
||
if d.DisplayName != "Дюна: Часть вторая (2024)" {
|
||
t.Errorf("display_name = %q", d.DisplayName)
|
||
}
|
||
// Извлечённая структура сохранена как parsed_context (базовый слой полей имени).
|
||
if d.ParsedContext == "" {
|
||
t.Error("parsed_context пуст, ожидалась сохранённая структура")
|
||
}
|
||
var pf naming.Fields
|
||
if err := json.Unmarshal([]byte(d.ParsedContext), &pf); err != nil {
|
||
t.Fatalf("parsed_context не JSON: %v", err)
|
||
}
|
||
if pf.Director != "Дени Вильнёв" {
|
||
t.Errorf("parsed_context.director = %q, want «Дени Вильнёв»", pf.Director)
|
||
}
|
||
}
|
||
|
||
// Торрент пойманной загрузки уже присутствует в qBittorrent (добавлен ранее
|
||
// вручную/другим клиентом) → усыновляем: promote в downloading без повторного
|
||
// Add (и без 409) и без namer; имя берём из раздачи снимка.
|
||
func TestProcessCatchedPresentAdopts(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "Дюна 2")
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{
|
||
{Hash: catchedIH, Name: "Dune.2024.1080p"}, // без нашей категории — добавлен вручную
|
||
}}
|
||
w := newTestWorker(st, qb)
|
||
nm := &fakeNamer{name: "не должно вызваться"}
|
||
w.SetNamer(nm)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("qbt.Add не должен вызываться для уже присутствующего торрента, calls = %d", len(qb.added))
|
||
}
|
||
if nm.calls != 0 {
|
||
t.Errorf("namer (LLM) не должен вызываться, calls = %d", nm.calls)
|
||
}
|
||
d := st.downloads["1"]
|
||
if d.State != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading", d.State)
|
||
}
|
||
if d.DisplayName != "Dune.2024.1080p" {
|
||
t.Errorf("display_name = %q, want имя раздачи из снимка", d.DisplayName)
|
||
}
|
||
}
|
||
|
||
// Гонка F6: список catched снят как magnet, но апгрейд до .torrent случился
|
||
// между снимком и re-read под замком. processCatched перечитывает source_type
|
||
// под w.mu, поэтому добавляет файлом (Torrents), а не magnet-ссылкой.
|
||
func TestProcessCatchedReReadsSourceTypeUnderLock(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx") // снят как magnet
|
||
st.torrents = map[string][]byte{}
|
||
qb := &fakeQbt{} // раздачи нет → absent-ветка (обычное добавление)
|
||
// Апгрейд «под носом»: между листингом (снимок magnet) и re-read под замком
|
||
// БД уже стала torrent с сохранёнными байтами.
|
||
qb.onTorrents = func() {
|
||
st.downloads["1"].SourceType = store.SourceTorrent
|
||
st.torrents["1"] = []byte("d4:infod-fake-torrent-bytes-ee")
|
||
}
|
||
w := newTestWorker(st, qb)
|
||
w.SetNamer(&fakeNamer{name: "X"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("qbt.Add calls = %d, want 1", len(qb.added))
|
||
}
|
||
if len(qb.added[0].Torrents) != 1 {
|
||
t.Errorf("после апгрейда ожидалось добавление файлом (Torrents), got Torrents=%v URLs=%v",
|
||
qb.added[0].Torrents, qb.added[0].URLs)
|
||
}
|
||
if len(qb.added[0].URLs) != 0 {
|
||
t.Errorf("magnet-ссылка не должна использоваться после апгрейда в torrent: %v", qb.added[0].URLs)
|
||
}
|
||
}
|
||
|
||
// Листинг qBittorrent провалился (недоступен) → пойманную не трогаем: остаётся
|
||
// catched, ни namer, ни Add не вызываются (повтор на следующем тике).
|
||
func TestProcessCatchedListErrorKeepsCatched(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{torrentsErr: errors.New("connection refused")}
|
||
w := newTestWorker(st, qb)
|
||
nm := &fakeNamer{name: "X"}
|
||
w.SetNamer(nm)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("state = %q, want catched (повтор)", st.downloads["1"].State)
|
||
}
|
||
if nm.calls != 0 || len(qb.added) != 0 {
|
||
t.Errorf("при недоступности qBit namer/Add не должны вызываться: namer=%d add=%d", nm.calls, len(qb.added))
|
||
}
|
||
}
|
||
|
||
// Транзиентный сбой add — остаёмся в catched для повтора на следующем тике.
|
||
func TestProcessCatchedTransientFailureKeepsCatched(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{addErr: errors.New("connection refused")}
|
||
w := newTestWorker(st, qb)
|
||
w.SetNamer(&fakeNamer{name: "X"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("state = %q, want catched (повтор)", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Предохранитель: catched старше catch_timeout → failed (qbit_add) + уведомление;
|
||
// add при этом не вызывается.
|
||
func TestProcessCatchedTimeoutFails(t *testing.T) {
|
||
old := store.FormatTime(time.Date(2026, 6, 14, 9, 0, 0, 0, time.UTC)) // 1 час до now
|
||
st := catchedStore("1", catchedIH, old, "ctx")
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
w.cfg.CatchTimeout = 10 * time.Minute
|
||
n := &recordingNotifier{ch: make(chan notifyEvent, 1)}
|
||
w.SetNotifier(n)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
d := st.downloads["1"]
|
||
if d.State != store.StateFailed || d.ErrorCode.String != errCodeQbitAdd {
|
||
t.Errorf("state = %q code = %q, want failed/qbit_add", d.State, d.ErrorCode.String)
|
||
}
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("add не должен вызываться при таймауте, calls = %d", len(qb.added))
|
||
}
|
||
if e := waitNotify(t, n); e.ev != EventFailed {
|
||
t.Errorf("событие = %q, want failed", e.ev)
|
||
}
|
||
}
|
||
|
||
// F3: отмена во время (медленного) вывода имени видна re-read'ом состояния ПЕРЕД
|
||
// Add — источник в qBittorrent не добавляется вовсе (раньше Add успевал пройти,
|
||
// оставляя неуправляемый торрент без владельца).
|
||
func TestProcessCatchedCancelledDuringNamerSkipsAdd(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
// namer имитирует параллельную отмену во время (медленного) вывода имени.
|
||
nm := &fakeNamer{name: "X", onCall: func() { st.downloads["1"].State = store.StateCancelled }}
|
||
w.SetNamer(nm)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("Add не должен вызываться после отмены (re-read перед add), calls = %d", len(qb.added))
|
||
}
|
||
if len(qb.deleted) != 0 {
|
||
t.Errorf("торрент не добавляли — удалять нечего, Delete calls = %d", len(qb.deleted))
|
||
}
|
||
if st.downloads["1"].State != store.StateCancelled {
|
||
t.Errorf("state = %q, want cancelled", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// F3, scoped cleanup: отмена приходит в окне ПОСЛЕ успешного Add, но до записи
|
||
// перехода (onAdd имитирует Cancel ровно между Add и PromoteCatched). Добавленный
|
||
// НАМИ торрент удаляется из qBittorrent С ДАННЫМИ (уборка своего артефакта).
|
||
func TestProcessCatchedCancelledAfterAddRemovesTorrent(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
qb.onAdd = func() { st.downloads["1"].State = store.StateCancelled }
|
||
w.SetNamer(&fakeNamer{name: "X"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("Add должен был вызваться, calls = %d", len(qb.added))
|
||
}
|
||
if len(qb.deleted) != 1 {
|
||
t.Fatalf("добавленный нами торрент должен быть удалён, Delete calls = %d", len(qb.deleted))
|
||
}
|
||
if !qb.deletedData[0] {
|
||
t.Error("уборка своего артефакта должна идти С ДАННЫМИ (deleteFiles=true)")
|
||
}
|
||
if len(qb.deleted[0]) == 0 || qb.deleted[0][0] != catchedIH {
|
||
t.Errorf("удаление не по infohash загрузки: %v", qb.deleted[0])
|
||
}
|
||
if st.downloads["1"].State != store.StateCancelled {
|
||
t.Errorf("state = %q, want cancelled (уборка не трогает состояние)", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// F3, негативный инвариант «удаляем только своё»: тот же infohash появился в
|
||
// qBittorrent во время namer (внешний клиент, окно гонки). Свежий листинг перед
|
||
// Add видит присутствие → Add не делаем И чужой торрент С ДАННЫМИ не удаляем.
|
||
func TestProcessCatchedPreexistingTorrentNotDeleted(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{} // снимок тика пуст → идём обычным путём добавления
|
||
w := newTestWorker(st, qb)
|
||
nm := &fakeNamer{name: "X", onCall: func() {
|
||
// внешний клиент добавил тот же торрент, пока выводилось имя
|
||
qb.torrents = []qbt.Torrent{{Hash: catchedIH, Name: "external"}}
|
||
}}
|
||
w.SetNamer(nm)
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("Add не должен вызываться: infohash уже присутствует перед add, calls = %d", len(qb.added))
|
||
}
|
||
if len(qb.deleted) != 0 {
|
||
t.Errorf("пред-существующий (чужой) торрент удалять нельзя, Delete calls = %d", len(qb.deleted))
|
||
}
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("state = %q, want catched (усыновление на следующем тике)", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// F3, safety-critical: свежий листинг присутствия ПЕРЕД add не удался (сеть
|
||
// отвалилась между тик-снимком и проверкой). Отсутствие infohash не подтверждено
|
||
// → Add не делаем (иначе delete-с-данными стал бы небезопасен), остаёмся в
|
||
// catched. onTorrents роняет ВТОРОЙ вызов Torrents (первый — тик-снимок).
|
||
func TestProcessCatchedPresenceRecheckFailKeepsCatched(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{}
|
||
calls := 0
|
||
qb.onTorrents = func() {
|
||
calls++
|
||
if calls == 2 { // тик-снимок (1) ок, свежий листинг перед add (2) падает
|
||
qb.torrentsErr = errors.New("connection refused")
|
||
}
|
||
}
|
||
w := newTestWorker(st, qb)
|
||
w.SetNamer(&fakeNamer{name: "X"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("Add не должен вызываться при неподтверждённом отсутствии, calls = %d", len(qb.added))
|
||
}
|
||
if len(qb.deleted) != 0 {
|
||
t.Errorf("Delete не должен вызываться, calls = %d", len(qb.deleted))
|
||
}
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("state = %q, want catched (повтор на следующем тике)", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// F3, различение «отмена vs сбой БД»: PromoteCatched упал транзиентно, но задача
|
||
// ЖИВА (state остался catched). Наш торрент НЕ удаляем — переход доведётся на
|
||
// следующем тике усыновлением присутствующей раздачи.
|
||
func TestProcessCatchedPromoteDBErrorKeepsTorrent(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
st.promoteErr = errors.New("db is locked") // сбой записи перехода, state = catched
|
||
qb := &fakeQbt{}
|
||
w := newTestWorker(st, qb)
|
||
w.SetNamer(&fakeNamer{name: "X"})
|
||
|
||
w.processCatched(context.Background())
|
||
|
||
if len(qb.added) != 1 {
|
||
t.Fatalf("Add должен был вызваться, calls = %d", len(qb.added))
|
||
}
|
||
if len(qb.deleted) != 0 {
|
||
t.Errorf("при транзиентном сбое БД (задача жива) торрент удалять нельзя, Delete calls = %d", len(qb.deleted))
|
||
}
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("state = %q, want catched (повтор промоушена на следующем тике)", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Поллинг активных (downloading) не трогает catched: раздачи в qBittorrent у
|
||
// пойманной загрузки ещё нет по дизайну, это не «пропажа».
|
||
func TestPollIgnoresCatched(t *testing.T) {
|
||
st := catchedStore("1", catchedIH, nowStr, "ctx")
|
||
qb := &fakeQbt{} // раздач нет
|
||
w := newTestWorker(st, qb)
|
||
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll: %v", err)
|
||
}
|
||
if st.downloads["1"].State != store.StateCatched {
|
||
t.Errorf("catched тронут поллингом: %q", st.downloads["1"].State)
|
||
}
|
||
}
|