Раздача активной (downloading) загрузки, исчезнувшая из qBittorrent (удалил пользователь/другой клиент), делала задачу вечным зомби: поллинг промахивался по torrentFor, писал Warn и continue каждый тик — состояние не менялось, уведомления и телеметрии не было, checkTimeouts без торрента не срабатывал. Пропажей источника у downloading не владел никто (сверка рассинхрона покрывает только done/target_missing/orphaned, восстановление — failed/stuck). Активный цикл Poll теперь применяет тот же дебаунс пропажи источника, что и сверка рассинхрона (source_miss_count / source_missing_threshold): после порога подряд идущих промахов задача уходит downloading → failed с distinct error_code source_gone и уведомлением. До порога транзиентная недоступность qBit (рестарт демона) задачу не роняет. source_gone восстановлению сверкой не подлежит (удаление намеренно), но штатно retriable — Retry заново отдаёт сохранённый источник; Retry сбрасывает source_miss_count, чтобы вернувшаяся задача получила полное грейс-окно, а не упала снова на ближайшем тике. Ребро downloading → failed уже было в графе, миграций/полей БД нет. Спека download-tracking дополнена требованием, диаграмма workflow.md — ребром. Change downloading-source-gone заархивирован. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
397 lines
19 KiB
Go
397 lines
19 KiB
Go
package worker
|
||
|
||
import (
|
||
"context"
|
||
"testing"
|
||
"time"
|
||
|
||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||
)
|
||
|
||
// addedRecent — added_on торрента «минуту назад» относительно зафиксированного
|
||
// в newTestWorker now (2026-06-14 10:00:00 UTC).
|
||
var addedRecent = time.Date(2026, 6, 14, 9, 59, 0, 0, time.UTC).Unix()
|
||
|
||
func oneFailed(state store.State, code, infohash, createdAt string) *fakeStore {
|
||
return &fakeStore{downloads: map[string]*store.Download{
|
||
"1": {
|
||
ID: "1",
|
||
State: state,
|
||
SourceType: store.SourceMagnet,
|
||
SourceRef: "magnet:?xt=urn:btih:" + infohash,
|
||
Infohashes: hashesOf("1", infohash),
|
||
ErrorCode: store.NullString(code),
|
||
CreatedAt: createdAt,
|
||
},
|
||
}}
|
||
}
|
||
|
||
func TestRecovery(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
tests := []struct {
|
||
name string
|
||
state store.State
|
||
code string
|
||
qbitState string
|
||
want store.State
|
||
}{
|
||
{"метаданные пришли → downloading", store.StateFailed, errCodeMagnetTimeout, "downloading", store.StateDownloading},
|
||
{"торрент готов → completed", store.StateFailed, errCodeMagnetTimeout, "uploading", store.StateCompleted},
|
||
{"всё ещё metaDL → остаётся failed", store.StateFailed, errCodeMagnetTimeout, "metaDL", store.StateFailed},
|
||
{"stalled ожил → downloading", store.StateStuck, errCodeStalled, "downloading", store.StateDownloading},
|
||
{"stalled всё ещё stalledDL → остаётся stuck", store.StateStuck, errCodeStalled, "stalledDL", store.StateStuck},
|
||
{"qbit_error не восстанавливается", store.StateFailed, errCodeQbitError, "downloading", store.StateFailed},
|
||
{"ошибка торрента не восстанавливает", store.StateFailed, errCodeMagnetTimeout, "error", store.StateFailed},
|
||
}
|
||
for _, tc := range tests {
|
||
t.Run(tc.name, func(t *testing.T) {
|
||
st := oneFailed(tc.state, tc.code, ih, timeOld)
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ih, State: tc.qbitState, AddedOn: addedRecent}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll: %v", err)
|
||
}
|
||
if got := st.downloads["1"].State; got != tc.want {
|
||
t.Errorf("state = %q, want %q", got, tc.want)
|
||
}
|
||
})
|
||
}
|
||
}
|
||
|
||
// Источник пропал (торрента нет в qBittorrent) — задача остаётся failed,
|
||
// воскрешать нечего (вернёт ручной retry).
|
||
func TestRecoveryNoSourceStaysFailed(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeMagnetTimeout, ih, timeOld)
|
||
w := newTestWorker(st, &fakeQbt{torrents: nil})
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if st.downloads["1"].State != store.StateFailed {
|
||
t.Errorf("без источника задача должна остаться failed, got %q", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Конфликт идемпотентности: тот же infohash уже взяла другая активная задача —
|
||
// упавшую не воскрешаем (иначе нарушим «одна активная задача на infohash»).
|
||
func TestRecoverySkipsOnIdempotencyConflict(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeMagnetTimeout, ih, timeOld)
|
||
st.downloads["2"] = &store.Download{
|
||
ID: "2",
|
||
State: store.StateDownloading,
|
||
SourceType: store.SourceMagnet,
|
||
Infohashes: hashesOf("2", ih),
|
||
CreatedAt: timeRecent,
|
||
}
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ih, State: "downloading", AddedOn: addedRecent}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if st.downloads["1"].State != store.StateFailed {
|
||
t.Errorf("при конфликте ключа задача #1 должна остаться failed, got %q", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Тот же конфликт ключа, но торрент уже готов (recovery хочет completed):
|
||
// completed тоже нетерминален и восстановил бы idempotency_key — проверка
|
||
// конфликта обязана покрывать и эту ветку.
|
||
func TestRecoverySkipsConflictOnCompleted(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeMagnetTimeout, ih, timeOld)
|
||
st.downloads["2"] = &store.Download{
|
||
ID: "2",
|
||
State: store.StateDownloading,
|
||
SourceType: store.SourceMagnet,
|
||
Infohashes: hashesOf("2", ih),
|
||
CreatedAt: timeRecent,
|
||
}
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ih, State: "uploading", AddedOn: addedRecent}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if st.downloads["1"].State != store.StateFailed {
|
||
t.Errorf("при конфликте ключа задача #1 не должна уходить в completed, got %q", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// Повторное падение одной задачи в пределах окна дебаунса шлёт уведомление лишь
|
||
// раз (защита от спама при флаппинге stuck↔downloading).
|
||
func TestFailNotifyDebounce(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateStuck, errCodeStalled, ih, timeOld)
|
||
w := newTestWorker(st, &fakeQbt{})
|
||
n := &recordingNotifier{ch: make(chan notifyEvent, 4)}
|
||
w.SetNotifier(n)
|
||
d := *st.downloads["1"]
|
||
|
||
w.transition(context.Background(), d, store.StateStuck, errCodeStalled, "")
|
||
if e := waitNotify(t, n); e.ev != EventFailed {
|
||
t.Fatalf("первый пинг: ev=%v, want failed", e.ev)
|
||
}
|
||
// Второе падение при том же w.now() — в пределах дебаунса, без пинга.
|
||
w.transition(context.Background(), d, store.StateStuck, errCodeStalled, "")
|
||
select {
|
||
case e := <-n.ch:
|
||
t.Fatalf("повторный пинг в пределах дебаунса не ожидался: %+v", e)
|
||
case <-time.After(200 * time.Millisecond):
|
||
}
|
||
}
|
||
|
||
// Retry при живом торренте перецепляется к нему (без повторного Add) и не падает
|
||
// снова на ближайшем тике: базис таймаута берётся от added_on, а не от старого
|
||
// created_at.
|
||
func TestRetryReattachesNoReadd(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeMagnetTimeout, ih, timeOld)
|
||
// Торрент жив, всё ещё тянет метаданные, но добавлен только что (added_on).
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ih, State: "metaDL", AddedOn: addedRecent}}}
|
||
w := newTestWorker(st, qb)
|
||
|
||
if err := w.Retry(context.Background(), "1"); err != nil {
|
||
t.Fatalf("Retry: %v", err)
|
||
}
|
||
if st.downloads["1"].State != store.StateDownloading {
|
||
t.Fatalf("после retry ожидался downloading, got %q", st.downloads["1"].State)
|
||
}
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("живой торрент не должен добавляться повторно, got %d Add", len(qb.added))
|
||
}
|
||
// Ближайший тик: metaDL свежий (added_on минуту назад) — не падает по таймауту.
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if st.downloads["1"].State != store.StateDownloading {
|
||
t.Errorf("свежий metaDL не должен падать после retry, got %q", st.downloads["1"].State)
|
||
}
|
||
}
|
||
|
||
// addedLongAgo — added_on «5 часов назад» относительно now теста (10:00:00 UTC):
|
||
// торрент давно в qBittorrent (возраст сам по себе большой).
|
||
var addedLongAgo = time.Date(2026, 6, 14, 5, 0, 0, 0, time.UTC).Unix()
|
||
|
||
// lastActivityLongAgo — last_activity «2 часа назад»: данные давно не двигались.
|
||
var lastActivityLongAgo = time.Date(2026, 6, 14, 8, 0, 0, 0, time.UTC).Unix()
|
||
|
||
// MAJOR-1: retry живого, но давно добавленного и простаивающего stalledDL-торрента
|
||
// сбрасывает базис таймаута (retried_at), поэтому СЛЕДУЮЩИЙ тик не роняет задачу
|
||
// снова в stuck. Без сброса базиса stallDuration=now−last_activity (2ч) > StuckAfter
|
||
// (1ч) → задача мгновенно вернулась бы в stuck (регрессия, которую прячет
|
||
// TestRetryReattachesNoReadd, ставящий added_on/last_activity «минуту назад»).
|
||
func TestRetryResetsTimeoutBasis(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateStuck, errCodeStalled, ih, timeOld)
|
||
// Торрент жив, но давно добавлен (added_on 5ч) и давно простаивает
|
||
// (last_activity 2ч) — по старой мере «возраст» он мгновенно снова stuck.
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{
|
||
Hash: ih, State: "stalledDL",
|
||
AddedOn: addedLongAgo, LastActivity: lastActivityLongAgo,
|
||
}}}
|
||
w := newTestWorker(st, qb)
|
||
|
||
if err := w.Retry(context.Background(), "1"); err != nil {
|
||
t.Fatalf("Retry: %v", err)
|
||
}
|
||
if len(qb.added) != 0 {
|
||
t.Errorf("живой здоровый торрент не должен добавляться повторно, got %d Add", len(qb.added))
|
||
}
|
||
if !st.downloads["1"].RetriedAt.Valid {
|
||
t.Error("retry должен проставить retried_at (сброс базиса)")
|
||
}
|
||
// Ключевая проверка MAJOR-1: следующий тик поллинга.
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateDownloading {
|
||
t.Errorf("после retry задача не должна снова падать в stuck на ближайшем тике, got %q", got)
|
||
}
|
||
}
|
||
|
||
// MAJOR-2: stuck_after мерит ДЛИТЕЛЬНОСТЬ ПРОСТОЯ (now−last_activity), а не возраст
|
||
// торрента. Долго качавшийся торрент (added_on 5ч назад) с недавним движением
|
||
// данных (last_activity 30с назад), на миг зашедший в stalledDL, НЕ уходит в stuck.
|
||
func TestStallMeasuredFromLastActivity(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
lastActivityRecent := time.Date(2026, 6, 14, 9, 59, 30, 0, time.UTC).Unix() // 30с назад
|
||
st := oneDownloading(ih, timeOld)
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{
|
||
Hash: ih, State: "stalledDL",
|
||
AddedOn: addedLongAgo, LastActivity: lastActivityRecent,
|
||
}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateDownloading {
|
||
t.Errorf("свежая активность (30с) — не stuck несмотря на возраст 5ч, got %q", got)
|
||
}
|
||
|
||
// Контроль: тот же торрент, но данные давно не двигались (last_activity 2ч) —
|
||
// простой превысил StuckAfter (1ч) → stuck.
|
||
qb.torrents[0].LastActivity = lastActivityLongAgo
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateStuck {
|
||
t.Errorf("простой 2ч > stuck_after 1ч должен дать stuck, got %q", got)
|
||
}
|
||
}
|
||
|
||
// Оборонительный клэмп: last_activity из будущего (перекос часов или sentinel
|
||
// «никогда не был активен») трактуется как непригодное значение. Иначе простой
|
||
// вышел бы отрицательным и реально застрявший торрент никогда не пометился бы
|
||
// stuck. Значение игнорируется, простой считается от базиса добавления (5ч).
|
||
func TestStallIgnoresFutureLastActivity(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
lastActivityFuture := time.Date(2026, 6, 14, 11, 0, 0, 0, time.UTC).Unix() // 1ч в будущем
|
||
st := oneDownloading(ih, timeOld)
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{
|
||
Hash: ih, State: "stalledDL",
|
||
AddedOn: addedLongAgo, LastActivity: lastActivityFuture,
|
||
}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateStuck {
|
||
t.Errorf("last_activity из будущего игнорируется (фолбэк на возраст 5ч > stuck_after 1ч → stuck), got %q", got)
|
||
}
|
||
}
|
||
|
||
// --- MAJOR-3: пропажа источника у активной (downloading) загрузки ---
|
||
|
||
// Источник активной загрузки устойчиво пропал из qBittorrent → задача уходит в
|
||
// failed(source_gone) с уведомлением, а не остаётся вечным зомби.
|
||
func TestSourceGoneMarksFailed(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneDownloading(ih, timeRecent)
|
||
w := newTestWorker(st, &fakeQbt{torrents: nil}) // раздачи нет
|
||
w.cfg.SourceMissingThreshold = 1 // помечаем при первой же пропаже
|
||
n := &recordingNotifier{ch: make(chan notifyEvent, 4)}
|
||
w.SetNotifier(n)
|
||
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll: %v", err)
|
||
}
|
||
d := st.downloads["1"]
|
||
if d.State != store.StateFailed {
|
||
t.Fatalf("state = %q, want failed", d.State)
|
||
}
|
||
if got := d.ErrorCode.String; got != errCodeSourceGone {
|
||
t.Errorf("error_code = %q, want %q", got, errCodeSourceGone)
|
||
}
|
||
if e := waitNotify(t, n); e.ev != EventFailed {
|
||
t.Errorf("пинг: ev=%v, want failed", e.ev)
|
||
}
|
||
}
|
||
|
||
// Кратковременная пропажа (меньше порога) не роняет задачу — дебаунс терпит
|
||
// транзиентную недоступность qBittorrent (рестарт демона).
|
||
func TestSourceGoneDebounced(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneDownloading(ih, timeRecent)
|
||
w := newTestWorker(st, &fakeQbt{torrents: nil})
|
||
w.cfg.SourceMissingThreshold = 3
|
||
|
||
for i := 1; i <= 2; i++ {
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll %d: %v", i, err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateDownloading {
|
||
t.Fatalf("tick %d: state = %q, want downloading (до порога)", i, got)
|
||
}
|
||
if got := st.downloads["1"].SourceMissCount; got != i {
|
||
t.Errorf("tick %d: miss = %d, want %d", i, got, i)
|
||
}
|
||
}
|
||
if err := w.Poll(context.Background()); err != nil { // третий промах
|
||
t.Fatalf("Poll 3: %v", err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateFailed {
|
||
t.Fatalf("tick 3: state = %q, want failed (порог достигнут)", got)
|
||
}
|
||
if got := st.downloads["1"].ErrorCode.String; got != errCodeSourceGone {
|
||
t.Errorf("error_code = %q, want %q", got, errCodeSourceGone)
|
||
}
|
||
}
|
||
|
||
// Возврат раздачи до порога сбрасывает счётчик промахов и задача ведётся
|
||
// обычной сверкой (не падает).
|
||
func TestSourceGoneResetOnReturn(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneDownloading(ih, timeRecent)
|
||
qb := &fakeQbt{torrents: nil}
|
||
w := newTestWorker(st, qb)
|
||
w.cfg.SourceMissingThreshold = 3
|
||
|
||
if err := w.Poll(context.Background()); err != nil { // промах 1
|
||
t.Fatalf("Poll miss: %v", err)
|
||
}
|
||
if got := st.downloads["1"].SourceMissCount; got != 1 {
|
||
t.Fatalf("miss = %d, want 1", got)
|
||
}
|
||
// Раздача вернулась (свежий metaDL — не падает по таймауту).
|
||
qb.torrents = []qbt.Torrent{{Hash: ih, State: "metaDL", AddedOn: addedRecent}}
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll return: %v", err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading (источник вернулся)", got)
|
||
}
|
||
if got := st.downloads["1"].SourceMissCount; got != 0 {
|
||
t.Errorf("miss = %d, want 0 (сброс)", got)
|
||
}
|
||
}
|
||
|
||
// source_gone восстановлению сверкой не подлежит: даже если раздача снова
|
||
// появилась и продвинулась, reconcileRecovery её не воскрешает — нужен ручной
|
||
// Retry.
|
||
func TestSourceGoneNotAutoRecovered(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeSourceGone, ih, timeOld)
|
||
qb := &fakeQbt{torrents: []qbt.Torrent{{Hash: ih, State: "downloading", AddedOn: addedRecent}}}
|
||
w := newTestWorker(st, qb)
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll: %v", err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateFailed {
|
||
t.Errorf("state = %q, want failed (source_gone сверкой не воскрешается)", got)
|
||
}
|
||
}
|
||
|
||
// Retry задачи source_gone сбрасывает source_miss_count, чтобы вернувшаяся в
|
||
// downloading задача получила полное грейс-окно, а не упала снова на ближайшем
|
||
// тике, если переотданная раздача ещё не видна в qBittorrent.
|
||
func TestRetryResetsSourceMissCount(t *testing.T) {
|
||
const ih = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
|
||
st := oneFailed(store.StateFailed, errCodeSourceGone, ih, timeRecent)
|
||
st.downloads["1"].SourceMissCount = 3 // задача упала по порогу
|
||
qb := &fakeQbt{torrents: nil} // раздачи всё ещё нет (Add вслепую)
|
||
w := newTestWorker(st, qb)
|
||
w.cfg.SourceMissingThreshold = 3
|
||
|
||
if err := w.Retry(context.Background(), "1"); err != nil {
|
||
t.Fatalf("Retry: %v", err)
|
||
}
|
||
if got := st.downloads["1"].SourceMissCount; got != 0 {
|
||
t.Fatalf("после retry miss = %d, want 0 (сброс)", got)
|
||
}
|
||
if len(qb.added) != 1 {
|
||
t.Errorf("retry без живой раздачи должен переотдать источник, got %d Add", len(qb.added))
|
||
}
|
||
// Грейс-окно: раздача ещё не видна — задача НЕ падает на ближайшем тике
|
||
// (miss стал бы 1 < 3, а не 4 при несброшенном счётчике).
|
||
if err := w.Poll(context.Background()); err != nil {
|
||
t.Fatalf("Poll after retry: %v", err)
|
||
}
|
||
if got := st.downloads["1"].State; got != store.StateDownloading {
|
||
t.Errorf("state = %q, want downloading (грейс-окно не съедено)", got)
|
||
}
|
||
if got := st.downloads["1"].SourceMissCount; got != 1 {
|
||
t.Errorf("miss = %d, want 1 (один промах после сброса)", got)
|
||
}
|
||
}
|