Восстановление зависших загрузок и уведомления о падении (state-reconciliation)
Долгий metaDL больше не убивается агрессивным таймаутом: дефолт magnet_timeout 30m → 24h (страховочный предохранитель), базис отсчёта — added_on из qBittorrent, а не created_at (переживает retry/усыновление). Авто-восстановление: фоновая сверка возвращает в поток задачи, упавшие по нашей нетерпеливости (magnet_timeout/stalled), когда источник ожил и продвинулся за условие падения (downloading/completed по статусу торрента); qbit_error не воскрешается. Конфликт idempotency (infohash занят другой активной задачей) — оставляем в failed. Уведомления: любой переход в failed/stuck пингует автора (включая приёмный qbit_add через ingest), с дебаунсом против спама при флаппинге stalled. Ручной retry добавлен в веб-UI и Telegram; Retry перецепляется к живому торренту вместо слепого Add. Дельта state-reconciliation влита в живые спеки; обновлён workflow.md. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -181,7 +181,7 @@ func Default() *Config {
|
||||
Worker: Worker{
|
||||
PollInterval: Duration(5 * time.Second),
|
||||
StuckAfter: Duration(time.Hour),
|
||||
MagnetTimeout: Duration(30 * time.Minute),
|
||||
MagnetTimeout: Duration(24 * time.Hour),
|
||||
SourceMissingThreshold: 3,
|
||||
},
|
||||
Recognition: Recognition{AutoConfidenceThreshold: 0.85},
|
||||
|
||||
@@ -84,6 +84,7 @@ func NewRouter(d Deps) (http.Handler, error) {
|
||||
r.Get("/", s.handleIndex)
|
||||
r.Post("/ui/downloads", s.handleUIAdd)
|
||||
r.Post("/ui/downloads/{id}/cancel", s.handleUICancel)
|
||||
r.Post("/ui/downloads/{id}/retry", s.handleUIRetry)
|
||||
|
||||
// Веб-UI: ревью раскладки.
|
||||
r.Get("/review/{id}", s.handleReview)
|
||||
@@ -133,6 +134,7 @@ type downloadView struct {
|
||||
Reviewable bool // review/deferred — есть экран ревью
|
||||
Undoable bool // done — можно откатить раскладку
|
||||
Relinkable bool // reverted/cancelled/target_missing — можно перепривязать заново
|
||||
Retriable bool // failed/stuck — можно повторить попытку
|
||||
Note string // пояснение рассинхрона (target_missing/orphaned/deleted)
|
||||
}
|
||||
|
||||
@@ -182,6 +184,19 @@ func (s *server) handleUICancel(w http.ResponseWriter, r *http.Request) {
|
||||
http.Redirect(w, r, "/", http.StatusSeeOther)
|
||||
}
|
||||
|
||||
func (s *server) handleUIRetry(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := pathID(r)
|
||||
if err != nil {
|
||||
redirectErr(w, r, "некорректный id")
|
||||
return
|
||||
}
|
||||
if err := s.deps.Commander.Retry(r.Context(), id); err != nil {
|
||||
redirectErr(w, r, userErr(r, err, id))
|
||||
return
|
||||
}
|
||||
http.Redirect(w, r, "/", http.StatusSeeOther)
|
||||
}
|
||||
|
||||
// --- REST API ---
|
||||
|
||||
type downloadDTO struct {
|
||||
@@ -320,7 +335,8 @@ func toView(d store.Download) downloadView {
|
||||
Undoable: d.State == store.StateDone,
|
||||
Relinkable: d.State == store.StateReverted || d.State == store.StateCancelled ||
|
||||
d.State == store.StateTargetMissing,
|
||||
Note: desyncNote(d.State),
|
||||
Retriable: d.State == store.StateFailed || d.State == store.StateStuck,
|
||||
Note: desyncNote(d.State),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -160,6 +160,23 @@ func TestAPICancel(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUIRetry(t *testing.T) {
|
||||
cmd := &fakeCommander{}
|
||||
srv := newServer(t, httpapi.Deps{Ingestor: &fakeIngestor{}, Commander: cmd, Reader: &fakeReader{}})
|
||||
|
||||
resp, err := http.Post(srv.URL+"/ui/downloads/5/retry", "application/x-www-form-urlencoded", nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK { // 303 → редирект на / → 200
|
||||
t.Fatalf("status = %d, want 200", resp.StatusCode)
|
||||
}
|
||||
if len(cmd.retried) != 1 || cmd.retried[0] != 5 {
|
||||
t.Errorf("retry вызван неверно: %v", cmd.retried)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAPICommandConflict(t *testing.T) {
|
||||
// Конфликт состояния (worker.ErrConflict) → 409, не 500.
|
||||
cmd := &fakeCommander{err: fmt.Errorf("cancel: download 5 in wrong state: %w", worker.ErrConflict)}
|
||||
|
||||
@@ -18,6 +18,10 @@ import (
|
||||
// capIngest — стадия приёма для поля capability в логах.
|
||||
const capIngest = "ingest"
|
||||
|
||||
// errCodeQbitAdd — error_code задачи, упавшей на добавлении источника в
|
||||
// qBittorrent (раздачи в qBittorrent нет, восстановлению не подлежит).
|
||||
const errCodeQbitAdd = "qbit_add"
|
||||
|
||||
// Store — нужная ingest часть хранилища.
|
||||
type Store interface {
|
||||
FindActiveByInfohash(ctx context.Context, infohash string) (*store.Download, error)
|
||||
@@ -49,6 +53,12 @@ type Service struct {
|
||||
namer Namer
|
||||
cfg Config
|
||||
log *slog.Logger
|
||||
|
||||
// notifyFailed — опц. пинг автору о падении приёма (добавление в qBittorrent
|
||||
// не удалось). Closure, а не worker.Notifier: приёмное падение в qBit не
|
||||
// попадает в поллинг-цикл worker (раздачи нет), поэтому уведомляет ingest
|
||||
// сам; closure избавляет ядро приёма от зависимости на пакет worker.
|
||||
notifyFailed func(downloadID int64)
|
||||
}
|
||||
|
||||
// New собирает сервис приёма. namer опционален (nil → отображаемое имя не
|
||||
@@ -57,6 +67,9 @@ func New(st Store, qb QBittorrent, namer Namer, cfg Config, log *slog.Logger) *S
|
||||
return &Service{store: st, qbt: qb, namer: namer, cfg: cfg, log: log}
|
||||
}
|
||||
|
||||
// SetFailureNotifier подключает пинг о падении приёма (до начала работы).
|
||||
func (s *Service) SetFailureNotifier(fn func(downloadID int64)) { s.notifyFailed = fn }
|
||||
|
||||
// Request — входной запрос приёма.
|
||||
type Request struct {
|
||||
Source string // пока — magnet-ссылка
|
||||
@@ -135,8 +148,12 @@ func (s *Service) Ingest(ctx context.Context, req Request) (Result, error) {
|
||||
// это разные факты, не дубль.
|
||||
log.Error("download accept failed", "error", addErr)
|
||||
// Задача уже в БД — помечаем failed, чтобы worker её не подхватил.
|
||||
if setErr := s.store.SetDownloadState(ctx, id, store.StateFailed, "qbit_add", addErr.Error()); setErr != nil {
|
||||
if setErr := s.store.SetDownloadState(ctx, id, store.StateFailed, errCodeQbitAdd, addErr.Error()); setErr != nil {
|
||||
log.Error("mark download failed after qbit error failed", "error", setErr)
|
||||
} else if s.notifyFailed != nil {
|
||||
// Это падение минует worker.transition (раздачи в qBit нет) — уведомляем
|
||||
// сами, чтобы приёмные провалы тоже доходили до автора.
|
||||
go s.notifyFailed(id)
|
||||
}
|
||||
return Result{DownloadID: id, Infohash: info.Infohash, State: store.StateFailed},
|
||||
fmt.Errorf("ingest: add to qbittorrent: %w", addErr)
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
@@ -177,6 +178,26 @@ func TestIngestQbitErrorMarksFailed(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestQbitErrorNotifies(t *testing.T) {
|
||||
fs := &fakeStore{}
|
||||
fq := &fakeQbt{err: errors.New("connection refused")}
|
||||
svc := newService(fs, fq)
|
||||
got := make(chan int64, 1)
|
||||
svc.SetFailureNotifier(func(id int64) { got <- id })
|
||||
|
||||
if _, err := svc.Ingest(context.Background(), Request{Source: sampleMagnet}); err == nil {
|
||||
t.Fatal("ожидалась ошибка")
|
||||
}
|
||||
select {
|
||||
case id := <-got:
|
||||
if id == 0 {
|
||||
t.Errorf("уведомление с нулевым id")
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("уведомление о падении приёма не пришло")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestRejectsNonMagnet(t *testing.T) {
|
||||
fs := &fakeStore{}
|
||||
fq := &fakeQbt{}
|
||||
|
||||
@@ -40,10 +40,13 @@ const (
|
||||
// ТОЛЬКО сюда — иначе семантика «активности» разъедется (idempotency_key
|
||||
// снимается по IsTerminal, а активность считалась бы по другому списку).
|
||||
//
|
||||
// Состояния рассинхрона (target_missing/orphaned/deleted) — терминальны по
|
||||
// тем же причинам, что reverted/cancelled: дальше двигает либо человек
|
||||
// (relink из target_missing), либо фоновая сверка (healing/прогрессия),
|
||||
// напрямую через SetDownloadState; ключ идемпотентности при этом не нужен.
|
||||
// Состояния рассинхрона (target_missing/orphaned/deleted), а также
|
||||
// failed/stuck — терминальны для idempotency_key, но не «мертвы»: дальше
|
||||
// двигает либо человек (relink из target_missing, retry из failed/stuck), либо
|
||||
// фоновая сверка (healing/прогрессия desync; авто-восстановление failed/stuck
|
||||
// при оживлении источника, см. state-reconciliation) — напрямую через
|
||||
// SetDownloadState, который восстановит ключ для нетерминального целевого
|
||||
// состояния.
|
||||
var terminalStates = []State{
|
||||
StateDone, StateCancelled, StateFailed, StateReverted,
|
||||
StateTargetMissing, StateOrphaned, StateDeleted,
|
||||
@@ -157,6 +160,30 @@ func (s *Store) ListDownloadsByState(ctx context.Context, states ...State) ([]Do
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// ListRecoverable возвращает задачи в failed/stuck с одним из переданных
|
||||
// error_code — кандидатов на авто-восстановление (см. state-reconciliation).
|
||||
// Фильтр по коду в SQL, чтобы не вычитывать на каждом тике поллинга все
|
||||
// накопленные провалы (qbit_error и пр.), которые восстановлению не подлежат.
|
||||
func (s *Store) ListRecoverable(ctx context.Context, codes ...string) ([]Download, error) {
|
||||
if len(codes) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
ph := make([]string, len(codes))
|
||||
args := make([]any, 0, len(codes)+2)
|
||||
for i, c := range codes {
|
||||
ph[i] = "?"
|
||||
args = append(args, c)
|
||||
}
|
||||
args = append(args, string(StateFailed), string(StateStuck))
|
||||
q := `SELECT * FROM download WHERE error_code IN (` + strings.Join(ph, ",") +
|
||||
`) AND state IN (?, ?) ORDER BY id DESC`
|
||||
var out []Download
|
||||
if err := s.DB.SelectContext(ctx, &out, q, args...); err != nil {
|
||||
return nil, fmt.Errorf("list recoverable: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// FindActiveByInfohash возвращает незавершённую задачу для infohash либо
|
||||
// (nil, nil), если её нет. Основа идемпотентного приёма.
|
||||
func (s *Store) FindActiveByInfohash(ctx context.Context, infohash string) (*Download, error) {
|
||||
|
||||
@@ -36,6 +36,7 @@ type Reviewer interface {
|
||||
SetType(ctx context.Context, id int64, mediaType string) error
|
||||
Defer(ctx context.Context, id int64) error
|
||||
Cancel(ctx context.Context, id int64) error
|
||||
Retry(ctx context.Context, id int64) error
|
||||
}
|
||||
|
||||
// Config — параметры бота.
|
||||
@@ -191,6 +192,9 @@ func (b *Bot) handleCallback(ctx context.Context, cq *tgbotapi.CallbackQuery) {
|
||||
case "reject":
|
||||
err = b.reviewer.Cancel(ctx, id)
|
||||
note = "Отклонено"
|
||||
case "retry":
|
||||
err = b.reviewer.Retry(ctx, id)
|
||||
note = "Повторяю…"
|
||||
case "type":
|
||||
err = b.reviewer.SetType(ctx, id, val)
|
||||
note = "Меняю тип…"
|
||||
@@ -248,6 +252,8 @@ func (b *Bot) Notify(ctx context.Context, downloadID int64, event worker.NotifyE
|
||||
text = b.renderDone(rd)
|
||||
case worker.EventTargetMissing, worker.EventOrphaned:
|
||||
text, kb = b.renderDesync(rd, event), b.webOnly(downloadID)
|
||||
case worker.EventFailed:
|
||||
text, kb = b.renderFailed(rd)
|
||||
default:
|
||||
text, kb = b.renderCard(rd)
|
||||
}
|
||||
|
||||
@@ -65,6 +65,7 @@ type fakeReviewer struct {
|
||||
typed map[int64]string
|
||||
deferred []int64
|
||||
canceled []int64
|
||||
retried []int64
|
||||
}
|
||||
|
||||
func (f *fakeReviewer) ReviewData(context.Context, int64) (*worker.ReviewData, error) {
|
||||
@@ -96,6 +97,10 @@ func (f *fakeReviewer) Cancel(_ context.Context, id int64) error {
|
||||
f.canceled = append(f.canceled, id)
|
||||
return nil
|
||||
}
|
||||
func (f *fakeReviewer) Retry(_ context.Context, id int64) error {
|
||||
f.retried = append(f.retried, id)
|
||||
return nil
|
||||
}
|
||||
|
||||
func reviewData(state store.State) *worker.ReviewData {
|
||||
s, e := 2, 1
|
||||
@@ -252,6 +257,29 @@ func TestBot_NotifyDone(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestBot_NotifyFailed(t *testing.T) {
|
||||
b, api, _, rev := newTestBot(t, []int64{7})
|
||||
rev.data = reviewData(store.StateFailed)
|
||||
b.Notify(context.Background(), 5, worker.EventFailed)
|
||||
|
||||
if len(api.sent) != 1 || !strings.Contains(api.sent[0].text, "не удалась") {
|
||||
t.Errorf("sent = %+v", api.sent)
|
||||
}
|
||||
if !api.sent[0].hasKB { // кнопка повтора
|
||||
t.Error("уведомление о падении без клавиатуры повтора")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBot_CallbackRetry(t *testing.T) {
|
||||
b, _, _, rev := newTestBot(t, []int64{7})
|
||||
rev.data = reviewData(store.StateFailed)
|
||||
b.handleCallback(context.Background(), cbFrom(7, "retry:5"))
|
||||
|
||||
if len(rev.retried) != 1 || rev.retried[0] != 5 {
|
||||
t.Errorf("retried = %v", rev.retried)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseCallback(t *testing.T) {
|
||||
a, id, v := parseCallback("type:5:series")
|
||||
if a != "type" || id != 5 || v != "series" {
|
||||
|
||||
@@ -31,6 +31,10 @@ func (b *Bot) renderCard(rd *worker.ReviewData) (string, *tgbotapi.InlineKeyboar
|
||||
if msg := rd.Download.ErrorMsg.String; msg != "" {
|
||||
text += "\n" + msg
|
||||
}
|
||||
// failed/stuck — даём кнопку повтора; остальное только «в вебе».
|
||||
if state == store.StateFailed || state == store.StateStuck {
|
||||
return text, b.retryKeyboard(id)
|
||||
}
|
||||
return text, b.webOnly(id)
|
||||
}
|
||||
}
|
||||
@@ -111,6 +115,41 @@ func (b *Bot) renderDesync(rd *worker.ReviewData, event worker.NotifyEvent) stri
|
||||
}
|
||||
}
|
||||
|
||||
// renderFailed — уведомление об упавшей/зависшей задаче с кнопкой повтора.
|
||||
func (b *Bot) renderFailed(rd *worker.ReviewData) (string, *tgbotapi.InlineKeyboardMarkup) {
|
||||
id := rd.Download.ID
|
||||
var sb strings.Builder
|
||||
verb := "не удалась"
|
||||
if rd.Download.State == store.StateStuck {
|
||||
verb = "зависла"
|
||||
}
|
||||
fmt.Fprintf(&sb, "❌ Задача #%d %s", id, verb)
|
||||
if code := rd.Download.ErrorCode.String; code != "" {
|
||||
fmt.Fprintf(&sb, " (%s)", code)
|
||||
}
|
||||
sb.WriteString(".")
|
||||
if msg := rd.Download.ErrorMsg.String; msg != "" {
|
||||
sb.WriteString("\n")
|
||||
sb.WriteString(msg)
|
||||
}
|
||||
if src := contextOrSource(rd); src != "" {
|
||||
fmt.Fprintf(&sb, "\nИсточник: %s", shorten(src, 80))
|
||||
}
|
||||
return sb.String(), b.retryKeyboard(id)
|
||||
}
|
||||
|
||||
// retryKeyboard — клавиатура для failed/stuck: повтор + опц. ссылка в веб.
|
||||
func (b *Bot) retryKeyboard(id int64) *tgbotapi.InlineKeyboardMarkup {
|
||||
row := []tgbotapi.InlineKeyboardButton{
|
||||
tgbotapi.NewInlineKeyboardButtonData("🔄 Повторить", "retry:"+itoa(id)),
|
||||
}
|
||||
if url := b.reviewURL(id); url != "" {
|
||||
row = append(row, tgbotapi.NewInlineKeyboardButtonURL("🌐 В вебе", url))
|
||||
}
|
||||
kb := tgbotapi.NewInlineKeyboardMarkup(tgbotapi.NewInlineKeyboardRow(row...))
|
||||
return &kb
|
||||
}
|
||||
|
||||
func (b *Bot) webOnly(id int64) *tgbotapi.InlineKeyboardMarkup {
|
||||
url := b.reviewURL(id)
|
||||
if url == "" {
|
||||
|
||||
@@ -160,6 +160,96 @@ func reconcileReason(sourcePresent, targetPresent bool) string {
|
||||
}
|
||||
}
|
||||
|
||||
// --- Восстановление зависших загрузок (failed/stuck → поток) ---
|
||||
|
||||
// reconcileRecovery воскрешает задачи, упавшие из-за нашей нетерпеливости
|
||||
// (magnet_timeout/stalled), когда их источник в qBittorrent ожил и продвинулся
|
||||
// за условие падения. Вызывается из Poll под w.mu. Реальные/пользовательские
|
||||
// провалы (qbit_error/reverted/cancelled/deleted) сюда не попадают.
|
||||
func (w *Worker) reconcileRecovery(ctx context.Context, byHash map[string]qbt.Torrent) {
|
||||
cands, err := w.store.ListRecoverable(ctx, errCodeMagnetTimeout, errCodeStalled)
|
||||
if err != nil {
|
||||
w.log.Warn("recovery list failed", "capability", capIngest, "error", err)
|
||||
return
|
||||
}
|
||||
for _, d := range cands {
|
||||
w.reconcileOneRecovery(ctx, d, byHash)
|
||||
}
|
||||
}
|
||||
|
||||
// reconcileOneRecovery возвращает одну зависшую задачу в поток, если её торрент
|
||||
// присутствует и продвинулся за условие падения.
|
||||
func (w *Worker) reconcileOneRecovery(ctx context.Context, d store.Download, byHash map[string]qbt.Torrent) {
|
||||
if !d.Infohash.Valid {
|
||||
return
|
||||
}
|
||||
t, ok := byHash[strings.ToLower(d.Infohash.String)]
|
||||
if !ok {
|
||||
return // источника нет — оставляем как есть (вернёт ручной retry)
|
||||
}
|
||||
if !torrentProgressed(d, t) {
|
||||
return // торрент всё ещё в metaDL/stalledDL или ошибочен — не воскрешаем
|
||||
}
|
||||
want := recoveredState(t.State)
|
||||
if want == "" {
|
||||
return // переходное состояние qBit (moving/checking) — ждём
|
||||
}
|
||||
ctx = w.scoped(ctx, capIngest, d.ID, d.Infohash.String)
|
||||
|
||||
// Конфликт идемпотентности: пока задача лежала в failed, тот же infohash мог
|
||||
// взять другая активная задача (idempotency_key снят при падении). Оба
|
||||
// целевых состояния (downloading/completed) нетерминальны → SetDownloadState
|
||||
// восстановит idempotency_key = infohash; при занятом ключе упёрлись бы в
|
||||
// unique-индекс. Поэтому проверяем владельца независимо от целевого состояния
|
||||
// и оставляем старую задачу в failed.
|
||||
other, err := w.store.FindActiveByInfohash(ctx, d.Infohash.String)
|
||||
if err != nil {
|
||||
logctx.From(ctx).Warn("recovery active lookup failed", "error", err)
|
||||
return
|
||||
}
|
||||
if other != nil && other.ID != d.ID {
|
||||
logctx.From(ctx).Info("recovery skipped, infohash taken by active download",
|
||||
"conflict_download_id", other.ID)
|
||||
return
|
||||
}
|
||||
// error_code/error_msg не пишем — задача снова здорова; причину в лог, а не в
|
||||
// поле ошибки (иначе она светилась бы в UI/REST как ошибка живой задачи).
|
||||
logctx.From(ctx).Info("recovery from failure", "to", want, "qbit_state", t.State)
|
||||
w.transition(ctx, d, want, "", "")
|
||||
}
|
||||
|
||||
// torrentProgressed сообщает, продвинулся ли торрент за условие, по которому
|
||||
// задача упала: для magnet_timeout — получил метаданные (вышел из metaDL); для
|
||||
// stalled — раздача ожила (вышла из stalledDL). Ошибочные состояния qBittorrent
|
||||
// продвижением не считаем (их ведёт обычный reconcile в qbit_error).
|
||||
func torrentProgressed(d store.Download, t qbt.Torrent) bool {
|
||||
if classify(t.State) == classErrored {
|
||||
return false
|
||||
}
|
||||
switch d.ErrorCode.String {
|
||||
case errCodeMagnetTimeout:
|
||||
return !isMeta(t.State)
|
||||
case errCodeStalled:
|
||||
return !isStalledDL(t.State)
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// recoveredState выводит состояние воскрешённой задачи из состояния торрента:
|
||||
// готов к раскладке → completed; ещё качается → downloading. Переходные
|
||||
// (moving/checking) и ошибочные состояния не восстанавливаем (пусто).
|
||||
func recoveredState(state string) store.State {
|
||||
switch classify(state) {
|
||||
case classReady:
|
||||
return store.StateCompleted
|
||||
case classDownloading:
|
||||
return store.StateDownloading
|
||||
default:
|
||||
return ""
|
||||
}
|
||||
}
|
||||
|
||||
// --- Синхронный preflight перед действием (не доверяем state в БД) ---
|
||||
|
||||
// ensureSourcePresent синхронно (без дебаунса) проверяет, что раздача есть в
|
||||
|
||||
@@ -0,0 +1,170 @@
|
||||
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[int64]*store.Download{
|
||||
1: {
|
||||
ID: 1,
|
||||
State: state,
|
||||
SourceType: store.SourceMagnet,
|
||||
SourceRef: "magnet:?xt=urn:btih:" + infohash,
|
||||
Infohash: store.NullString(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,
|
||||
Infohash: store.NullString(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,
|
||||
Infohash: store.NullString(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)
|
||||
}
|
||||
}
|
||||
@@ -250,6 +250,22 @@ func (m *memStore) ListDownloadsByState(_ context.Context, states ...store.State
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (m *memStore) ListRecoverable(_ context.Context, codes ...string) ([]store.Download, error) {
|
||||
var out []store.Download
|
||||
for _, d := range m.downloads {
|
||||
if d.State != store.StateFailed && d.State != store.StateStuck {
|
||||
continue
|
||||
}
|
||||
for _, c := range codes {
|
||||
if d.ErrorCode.Valid && d.ErrorCode.String == c {
|
||||
out = append(out, *d)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (m *memStore) ExistsByInfohash(_ context.Context, infohash string) (bool, error) {
|
||||
for _, d := range m.downloads {
|
||||
if d.Infohash.Valid && d.Infohash.String == infohash {
|
||||
|
||||
+98
-20
@@ -38,6 +38,7 @@ const (
|
||||
// Store — нужная worker часть хранилища.
|
||||
type Store interface {
|
||||
ListDownloadsByState(ctx context.Context, states ...store.State) ([]store.Download, error)
|
||||
ListRecoverable(ctx context.Context, codes ...string) ([]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
|
||||
@@ -94,6 +95,17 @@ const (
|
||||
EventDone NotifyEvent = "done" // раскладка завершена
|
||||
EventOrphaned NotifyEvent = "orphaned" // источник пропал, цель — последняя копия
|
||||
EventTargetMissing NotifyEvent = "target_missing" // цель удалена, доступен relink
|
||||
EventFailed NotifyEvent = "failed" // задача упала/зависла (failed/stuck)
|
||||
)
|
||||
|
||||
// Коды ошибок (error_code) при переходе в failed/stuck. Восстановимые
|
||||
// (magnet_timeout/stalled) — следствие нашей нетерпеливости: сверка воскрешает
|
||||
// такие задачи при оживлении источника (см. reconcileRecovery). qbit_error —
|
||||
// реальная ошибка qBittorrent, восстановлению не подлежит.
|
||||
const (
|
||||
errCodeMagnetTimeout = "magnet_timeout"
|
||||
errCodeStalled = "stalled"
|
||||
errCodeQbitError = "qbit_error"
|
||||
)
|
||||
|
||||
// Notifier — исходящие пинги (Telegram). Вызывается неблокирующе.
|
||||
@@ -136,8 +148,18 @@ type Worker struct {
|
||||
newID func() string // генератор apply_batch_id (подменяется в тестах)
|
||||
notifier Notifier // опц. исходящие пинги
|
||||
scanner Scanner // опц. пересканирование Jellyfin
|
||||
|
||||
// failNotified — дебаунс повторных EventFailed по задаче (download_id →
|
||||
// время последнего пинга). Мерцающий stalled-торрент колеблется
|
||||
// stuck↔downloading; без дебаунса каждый цикл слал бы уведомление. Память
|
||||
// процесса: при рестарте дебаунс сбрасывается — допустимо. Доступ под w.mu.
|
||||
failNotified map[int64]time.Time
|
||||
}
|
||||
|
||||
// failNotifyDebounce — минимальный интервал между уведомлениями о падении
|
||||
// одной задачи (см. failNotified).
|
||||
const failNotifyDebounce = time.Hour
|
||||
|
||||
// SetNotifier подключает исходящие пинги (до запуска Run).
|
||||
func (w *Worker) SetNotifier(n Notifier) { w.notifier = n }
|
||||
|
||||
@@ -148,14 +170,15 @@ func (w *Worker) SetScanner(s Scanner) { w.scanner = s }
|
||||
// распознавания и раскладки) — тогда completed-задачи не двигаются дальше.
|
||||
func New(st Store, qb QBittorrent, rec Recognizer, lay Layouter, cfg Config, log *slog.Logger) *Worker {
|
||||
return &Worker{
|
||||
store: st,
|
||||
qbt: qb,
|
||||
recognizer: rec,
|
||||
layouter: lay,
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
now: time.Now,
|
||||
newID: defaultBatchID,
|
||||
store: st,
|
||||
qbt: qb,
|
||||
recognizer: rec,
|
||||
layouter: lay,
|
||||
cfg: cfg,
|
||||
log: log,
|
||||
now: time.Now,
|
||||
newID: defaultBatchID,
|
||||
failNotified: map[int64]time.Time{},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -246,6 +269,10 @@ func (w *Worker) Poll(ctx context.Context) error {
|
||||
// Сверка разложенных задач с реальностью (источник в qBit + хардлинки на ФС)
|
||||
// — отдельно от активных, по двумерной матрице (см. state-reconciliation).
|
||||
w.reconcileDesync(ctx, byHash)
|
||||
|
||||
// Восстановление задач, упавших по нашей нетерпеливости (magnet_timeout/
|
||||
// stalled), если их источник в qBittorrent ожил и продвинулся.
|
||||
w.reconcileRecovery(ctx, byHash)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -257,7 +284,7 @@ func (w *Worker) reconcile(ctx context.Context, d store.Download, t qbt.Torrent)
|
||||
case classReady:
|
||||
w.transition(ctx, d, store.StateCompleted, "", "")
|
||||
case classErrored:
|
||||
w.transition(ctx, d, store.StateFailed, "qbit_error", "qBittorrent state: "+t.State)
|
||||
w.transition(ctx, d, store.StateFailed, errCodeQbitError, "qBittorrent state: "+t.State)
|
||||
case classDownloading:
|
||||
w.checkTimeouts(ctx, d, t)
|
||||
case classBusy:
|
||||
@@ -265,26 +292,43 @@ func (w *Worker) reconcile(ctx context.Context, d store.Download, t qbt.Torrent)
|
||||
}
|
||||
}
|
||||
|
||||
// checkTimeouts помечает зависшие задачи. Возраст считаем от created_at:
|
||||
// для metaDL это время с момента добавления (огрублённо, но достаточно).
|
||||
// checkTimeouts помечает зависшие задачи. Возраст считаем от факта в
|
||||
// qBittorrent (added_on), а не от created_at: базис переживает retry и
|
||||
// усыновление раздачи (см. design download-failure-recovery). magnet_timeout —
|
||||
// редкий страховочный предохранитель (дефолт 24h); настоящие провалы ловит
|
||||
// classErrored, а ожившие задачи воскрешает reconcileRecovery.
|
||||
func (w *Worker) checkTimeouts(ctx context.Context, d store.Download, t qbt.Torrent) {
|
||||
created, err := d.CreatedTime()
|
||||
if err != nil {
|
||||
logctx.From(ctx).Warn("cannot parse created_at", "value", d.CreatedAt, "error", err)
|
||||
return
|
||||
}
|
||||
age := w.now().Sub(created)
|
||||
age := w.torrentAge(d, t)
|
||||
|
||||
switch {
|
||||
case isMeta(t.State) && w.cfg.MagnetTimeout > 0 && age > w.cfg.MagnetTimeout:
|
||||
w.transition(ctx, d, store.StateFailed, "magnet_timeout",
|
||||
w.transition(ctx, d, store.StateFailed, errCodeMagnetTimeout,
|
||||
fmt.Sprintf("no metadata after %s", age.Truncate(time.Second)))
|
||||
case isStalledDL(t.State) && w.cfg.StuckAfter > 0 && age > w.cfg.StuckAfter:
|
||||
w.transition(ctx, d, store.StateStuck, "stalled",
|
||||
w.transition(ctx, d, store.StateStuck, errCodeStalled,
|
||||
fmt.Sprintf("stalled for %s", age.Truncate(time.Second)))
|
||||
}
|
||||
}
|
||||
|
||||
// torrentAge — возраст торрента: от added_on в qBittorrent (надёжный базис,
|
||||
// переживает retry/усыновление), с фолбэком на created_at задачи, если qBit не
|
||||
// отдал added_on.
|
||||
func (w *Worker) torrentAge(d store.Download, t qbt.Torrent) time.Duration {
|
||||
if t.AddedOn > 0 {
|
||||
return w.now().Sub(time.Unix(t.AddedOn, 0).UTC())
|
||||
}
|
||||
created, err := d.CreatedTime()
|
||||
if err != nil {
|
||||
// Ни added_on от qBit, ни разбираемого created_at — возраст неизвестен,
|
||||
// таймауты не сработают; фиксируем диагностикой.
|
||||
w.log.Warn("cannot determine torrent age",
|
||||
"capability", capIngest, "download_id", d.ID,
|
||||
"created_at", d.CreatedAt, "error", err)
|
||||
return 0
|
||||
}
|
||||
return w.now().Sub(created)
|
||||
}
|
||||
|
||||
// transition пишет новое состояние и логирует переход.
|
||||
func (w *Worker) transition(ctx context.Context, d store.Download, state store.State, code, msg string) {
|
||||
// FromOr, а не From: если вызывающий не завёл scoped-логгер, падаем на
|
||||
@@ -308,6 +352,10 @@ func (w *Worker) transition(ctx context.Context, d store.Download, state store.S
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventOrphaned)
|
||||
case store.StateTargetMissing:
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventTargetMissing)
|
||||
case store.StateFailed, store.StateStuck:
|
||||
if w.shouldNotifyFail(d.ID) {
|
||||
go w.notifier.Notify(context.Background(), d.ID, EventFailed)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -323,6 +371,25 @@ func (w *Worker) transition(ctx context.Context, d store.Download, state store.S
|
||||
}
|
||||
}
|
||||
|
||||
// shouldNotifyFail дебаунсит повторные уведомления о падении одной задачи
|
||||
// (мерцающий stalled-торрент: stuck↔downloading), чтобы не спамить. Вызывается
|
||||
// под w.mu. НЕ сбрасываем запись при восстановлении — иначе дебаунс не гасил бы
|
||||
// флаппинг.
|
||||
func (w *Worker) shouldNotifyFail(id int64) bool {
|
||||
now := w.now()
|
||||
if last, ok := w.failNotified[id]; ok && now.Sub(last) < failNotifyDebounce {
|
||||
return false
|
||||
}
|
||||
w.failNotified[id] = now
|
||||
// Лёгкая чистка устаревших записей, чтобы карта не росла без предела.
|
||||
for k, t := range w.failNotified {
|
||||
if now.Sub(t) >= failNotifyDebounce {
|
||||
delete(w.failNotified, k)
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// Cancel отклоняет задачу. Торрент в qBittorrent не трогаем — он продолжает
|
||||
// раздачу (источник неприкосновенен).
|
||||
func (w *Worker) Cancel(ctx context.Context, id int64) error {
|
||||
@@ -356,7 +423,18 @@ func (w *Worker) Retry(ctx context.Context, id int64) error {
|
||||
if d.State != store.StateFailed && d.State != store.StateStuck {
|
||||
return fmt.Errorf("retry: download %d is %s, only failed/stuck are retriable", id, d.State)
|
||||
}
|
||||
if d.SourceType == store.SourceMagnet {
|
||||
// Если раздача уже жива в qBittorrent — перецепляемся к ней, повторный Add
|
||||
// не нужен (и вреден: вслепую дублировал бы торрент). Add — только когда
|
||||
// источника в qBittorrent нет. Базис таймаута берётся от added_on, поэтому
|
||||
// возврат в downloading не роняет задачу снова на ближайшем тике.
|
||||
alive := false
|
||||
if d.Infohash.Valid {
|
||||
_, alive, err = w.torrentByInfohash(ctx, d.Infohash.String)
|
||||
if err != nil {
|
||||
return fmt.Errorf("retry: %w", err)
|
||||
}
|
||||
}
|
||||
if !alive && d.SourceType == store.SourceMagnet {
|
||||
if err := w.qbt.Add(ctx, qbt.AddRequest{
|
||||
URLs: []string{d.SourceRef},
|
||||
Category: w.cfg.Category,
|
||||
|
||||
@@ -42,6 +42,22 @@ func (f *fakeStore) ListDownloadsByState(_ context.Context, states ...store.Stat
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (f *fakeStore) ListRecoverable(_ context.Context, codes ...string) ([]store.Download, error) {
|
||||
var out []store.Download
|
||||
for _, d := range f.downloads {
|
||||
if d.State != store.StateFailed && d.State != store.StateStuck {
|
||||
continue
|
||||
}
|
||||
for _, c := range codes {
|
||||
if d.ErrorCode.Valid && d.ErrorCode.String == c {
|
||||
out = append(out, *d)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (f *fakeStore) GetDownload(_ context.Context, id int64) (*store.Download, error) {
|
||||
d, ok := f.downloads[id]
|
||||
if !ok {
|
||||
|
||||
Reference in New Issue
Block a user