Retry/stall: сброс базиса таймаута + простой от last_activity (MAJOR-1, MAJOR-2)

Два связанных бага семантики таймаутов зависания и ручного retry.

MAJOR-1: Retry живого торрента не сбрасывал базис отсчёта таймаута — задача
мгновенно снова падала в stuck на ближайшем тике. Вводим колонку
download.retried_at (миграция 0010): ручной retry фиксирует момент и
приподнимает пол обоих таймаутов (max(базис, retried_at)). Хранится в БД, а
не в памяти, чтобы сброс пережил тик поллинга и рестарт.

MAJOR-2: stuck_after мерил ВОЗРАСТ торрента (от added_on), а не ПРОСТОЙ —
долго качавшийся торрент, на миг зашедший в stalledDL, ложно уходил в stuck
со «stalled for 5h». Теперь stuck_after мерит простой от qBit last_activity
(новое поле qbt.Torrent из того же ответа /torrents/info); magnet_timeout
по-прежнему мерит возраст (семантически верно). checkTimeouts разбит на
torrentAge/stallDuration/addedBasis/retriedFloor.

NIT-10: фолбэк базиса возраста added_on→created_at сохранён и покрыт.
NIT-12: retry перестаёт перецепляться к сломанному живому торренту
(error/missingFiles) — повторно отдаёт источник (перецепка к нему
бессмысленна: reconcile тут же вернул бы в failed).

Спека: дельта state-reconciliation (MODIFIED «Восстановление зависшей
загрузки» и «Ручной повтор»), правка docs/specs/workflow.md (устранено
противоречие «возраст vs простой»), ER-схема database.md.

Тесты: TestRetryResetsTimeoutBasis (следующий тик после retry — прячется в
TestRetryReattachesNoReadd), TestStallMeasuredFromLastActivity,
TestSetRetriedAtOverwrites.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
av
2026-07-08 17:18:28 +03:00
co-authored by Claude Opus 4.8
parent 4475fbd548
commit 8261d5b55d
15 changed files with 659 additions and 40 deletions
+7 -2
View File
@@ -60,8 +60,13 @@ type Torrent struct {
AmountLeft int64 `json:"amount_left"`
TotalSize int64 `json:"total_size"` // полный размер раздачи, байт
AddedOn int64 `json:"added_on"`
InfohashV1 string `json:"infohash_v1"`
InfohashV2 string `json:"infohash_v2"`
// LastActivity — Unix-время последнего движения данных по торренту (скачан/
// отдан кусок). Базис измерения простоя для stuck_after: простой = now
// last_activity (а не возраст от added_on), поэтому долго качавшийся торрент,
// на миг зашедший в stalledDL, не помечается «зависшим».
LastActivity int64 `json:"last_activity"`
InfohashV1 string `json:"infohash_v1"`
InfohashV2 string `json:"infohash_v2"`
// Живая телеметрия (для снимка воркера и веб-UI).
Dlspeed int64 `json:"dlspeed"` // скорость загрузки, байт/с
+36 -2
View File
@@ -183,8 +183,13 @@ type Download struct {
// сортировки списка. NULL, пока воркер не наблюдал раздачу. Хранится в
// формате RFC 3339 (UTC, суффикс Z), как created_at.
SourceAddedAt sql.NullString `db:"source_added_at"`
CreatedAt string `db:"created_at"`
UpdatedAt string `db:"updated_at"`
// RetriedAt — время последнего ручного retry (RFC 3339 UTC, суффикс Z), NULL
// пока задачу не повторяли. Приподнимает базис отсчёта таймаутов, чтобы
// возврат в downloading не ронял задачу снова на ближайшем тике (см.
// state-reconciliation «Ручной повтор»).
RetriedAt sql.NullString `db:"retried_at"`
CreatedAt string `db:"created_at"`
UpdatedAt string `db:"updated_at"`
// Infohashes — хеши загрузки (download_infohash); подгружаются вместе с
// записью методами чтения store (v1 раньше v2 — порядок стабильный).
@@ -236,6 +241,19 @@ func Now() time.Time { return time.Now().UTC() }
// CreatedTime возвращает время создания загрузки как time.Time (UTC).
func (d Download) CreatedTime() (time.Time, error) { return ParseTime(d.CreatedAt) }
// RetriedTime возвращает время последнего ручного retry (UTC) и ok=false, если
// задачу ещё не повторяли (retried_at NULL) или метку не разобрать.
func (d Download) RetriedTime() (time.Time, bool) {
if !d.RetriedAt.Valid {
return time.Time{}, false
}
t, err := ParseTime(d.RetriedAt.String)
if err != nil {
return time.Time{}, false
}
return t, true
}
// NullString строит sql.NullString: пустая строка → NULL.
func NullString(s string) sql.NullString {
return sql.NullString{String: s, Valid: s != ""}
@@ -504,6 +522,22 @@ func (s *Store) SetSourceAddedAt(ctx context.Context, id string, t time.Time) er
return nil
}
// SetRetriedAt проставляет время ручного retry задачи (сброс базиса отсчёта
// таймаутов, см. state-reconciliation «Ручной повтор»). В отличие от
// SetSourceAddedAt перезаписывает значение: retry можно повторять, и базис
// должен смещаться на каждый.
func (s *Store) SetRetriedAt(ctx context.Context, id string, t time.Time) error {
res, err := s.DB.ExecContext(ctx,
`UPDATE download SET retried_at = ? WHERE id = ?`, FormatTime(t), id)
if err != nil {
return fmt.Errorf("set download %s retried_at: %w", id, err)
}
if n, _ := res.RowsAffected(); n == 0 {
return fmt.Errorf("set download %s retried_at: not found", id)
}
return nil
}
// ListDownloads возвращает все загрузки, новые сверху (id — ULID, сортировка
// по нему хронологична).
func (s *Store) ListDownloads(ctx context.Context) ([]Download, error) {
+28
View File
@@ -225,6 +225,34 @@ func TestSetSourceAddedAtOnce(t *testing.T) {
}
}
func TestSetRetriedAtOverwrites(t *testing.T) {
st := newTestStore(t)
ctx := context.Background()
id := mkDownload(t, st, 1, StateStuck, "x")
first := time.Date(2026, 6, 1, 10, 0, 0, 0, time.UTC)
second := time.Date(2026, 6, 2, 10, 0, 0, 0, time.UTC)
if err := st.SetRetriedAt(ctx, id, first); err != nil {
t.Fatal(err)
}
// В отличие от source_added_at — повторный retry перезаписывает базис.
if err := st.SetRetriedAt(ctx, id, second); err != nil {
t.Fatal(err)
}
d, err := st.GetDownload(ctx, id)
if err != nil {
t.Fatal(err)
}
if !d.RetriedAt.Valid || d.RetriedAt.String != FormatTime(second) {
t.Fatalf("retried_at = %q, want %q (последний retry)", d.RetriedAt.String, FormatTime(second))
}
got, ok := d.RetriedTime()
if !ok || !got.Equal(second) {
t.Fatalf("RetriedTime() = %v, %v; want %v, true", got, ok, second)
}
}
func TestListDownloadsPageRecTitleFallback(t *testing.T) {
st := newTestStore(t)
ctx := context.Background()
@@ -0,0 +1,11 @@
-- +goose Up
-- Время последнего ручного retry задачи (RFC 3339 UTC, суффикс Z). Приподнимает
-- базис отсчёта таймаутов (magnet_timeout/stuck_after): после retry отсчёт идёт
-- от max(базис_добавления или last_activity, retried_at), чтобы возврат в
-- downloading не ронял задачу снова на ближайшем тике (см. state-reconciliation
-- «Ручной повтор зависшей/упавшей загрузки»). Хранится в БД, а не в памяти,
-- чтобы сброс базиса пережил интервал поллинга и рестарт процесса.
ALTER TABLE download ADD COLUMN retried_at TEXT;
-- +goose Down
ALTER TABLE download DROP COLUMN retried_at;
+71
View File
@@ -168,3 +168,74 @@ func TestRetryReattachesNoReadd(t *testing.T) {
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=nowlast_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 мерит ДЛИТЕЛЬНОСТЬ ПРОСТОЯ (nowlast_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)
}
}
+7
View File
@@ -467,6 +467,13 @@ func (m *memStore) SetSourceAddedAt(_ context.Context, id string, t time.Time) e
return nil
}
func (m *memStore) SetRetriedAt(_ context.Context, id string, t time.Time) error {
if d, ok := m.downloads[id]; ok {
d.RetriedAt = store.NullString(store.FormatTime(t))
}
return nil
}
func (m *memStore) CreateRecognition(_ context.Context, r *store.Recognition, reasons []string) (string, error) {
for _, e := range m.recs {
if e.DownloadID == r.DownloadID {
+92 -28
View File
@@ -51,6 +51,10 @@ type Store interface {
PromoteCatched(ctx context.Context, id, displayName string) error
SetSourceMissCount(ctx context.Context, id string, n int) error
SetSourceAddedAt(ctx context.Context, id string, t time.Time) error
// SetRetriedAt проставляет время ручного retry — сброс базиса отсчёта
// таймаутов (magnet_timeout/stuck_after), чтобы возврат в downloading не
// ронял задачу снова на ближайшем тике.
SetRetriedAt(ctx context.Context, id string, t time.Time) error
// Идентичность/инвариант «одна активная загрузка на infohash».
ExistsByInfohash(ctx context.Context, hashes ...string) (bool, error)
@@ -506,21 +510,29 @@ func (w *Worker) reconcile(ctx context.Context, d store.Download, t qbt.Torrent)
}
}
// checkTimeouts помечает зависшие задачи. Возраст считаем от факта в
// qBittorrent (added_on), а не от created_at: базис переживает retry и
// усыновление раздачи (см. design download-failure-recovery). magnet_timeout —
// редкий страховочный предохранитель (дефолт 24h); настоящие провалы ловит
// checkTimeouts помечает зависшие задачи двумя разными мерами:
// - magnet_timeout — по ВОЗРАСТУ торрента (metaDL дольше magnet_timeout без
// метаданных): страховочный предохранитель (дефолт 24h), базис — added_on;
// - stuck_after — по ДЛИТЕЛЬНОСТИ ПРОСТОЯ (stalledDL без движения данных
// дольше stuck_after): базис — last_activity, а не возраст. Иначе долго
// качавшийся торрент, на миг зашедший в stalledDL, ложно уходит в stuck со
// «stalled for 5h» (см. state-reconciliation, MAJOR-2).
//
// Оба базиса приподняты до retried_at (ручной retry), чтобы возврат в
// downloading не ронял задачу снова на ближайшем тике. Настоящие провалы ловит
// classErrored, а ожившие задачи воскрешает reconcileRecovery.
func (w *Worker) checkTimeouts(ctx context.Context, d store.Download, t qbt.Torrent) {
age := w.torrentAge(d, t)
switch {
case isMeta(t.State) && w.cfg.MagnetTimeout > 0 && age > w.cfg.MagnetTimeout:
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, errCodeStalled,
fmt.Sprintf("stalled for %s", age.Truncate(time.Second)))
case isMeta(t.State) && w.cfg.MagnetTimeout > 0:
if age, ok := w.torrentAge(d, t); ok && age > w.cfg.MagnetTimeout {
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:
if idle, ok := w.stallDuration(d, t); ok && idle > w.cfg.StuckAfter {
w.transition(ctx, d, store.StateStuck, errCodeStalled,
fmt.Sprintf("stalled for %s", idle.Truncate(time.Second)))
}
}
}
@@ -572,23 +584,61 @@ func (w *Worker) captureInfohashes(ctx context.Context, d store.Download, t qbt.
}
}
// torrentAge — возраст торрента: от added_on в qBittorrent (надёжный базис,
// переживает retry/усыновление), с фолбэком на created_at задачи, если qBit не
// отдал added_on.
func (w *Worker) torrentAge(d store.Download, t qbt.Torrent) time.Duration {
// torrentAge — возраст торрента для magnet_timeout: от added_on в qBittorrent
// (надёжный базис, переживает усыновление), с фолбэком на created_at задачи,
// если qBit не отдал added_on (NIT-10). Базис приподнят до retried_at, чтобы
// ручной retry сбрасывал отсчёт. ok=false — базис неизвестен (ни added_on, ни
// разбираемого created_at): таймаут не срабатывает, фиксируем диагностикой.
func (w *Worker) torrentAge(d store.Download, t qbt.Torrent) (time.Duration, bool) {
basis, ok := w.addedBasis(d, t)
if !ok {
return 0, false
}
return w.now().Sub(w.retriedFloor(d, basis)), true
}
// stallDuration — длительность простоя торрента для stuck_after: от
// last_activity qBittorrent (момент последнего движения данных), с фолбэком на
// базис добавления, если qBit не отдал last_activity. Базис приподнят до
// retried_at (ручной retry даёт свежее окно). ok=false — базис неизвестен.
func (w *Worker) stallDuration(d store.Download, t qbt.Torrent) (time.Duration, bool) {
var basis time.Time
if t.LastActivity > 0 {
basis = time.Unix(t.LastActivity, 0).UTC()
} else {
var ok bool
if basis, ok = w.addedBasis(d, t); !ok {
return 0, false
}
}
return w.now().Sub(w.retriedFloor(d, basis)), true
}
// addedBasis — момент добавления торрента: added_on qBittorrent, иначе
// created_at задачи (NIT-10). ok=false — ни того, ни другого разобрать не
// удалось; фиксируем диагностикой.
func (w *Worker) addedBasis(d store.Download, t qbt.Torrent) (time.Time, bool) {
if t.AddedOn > 0 {
return w.now().Sub(time.Unix(t.AddedOn, 0).UTC())
return time.Unix(t.AddedOn, 0).UTC(), true
}
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 time.Time{}, false
}
return w.now().Sub(created)
return created, true
}
// retriedFloor приподнимает базис отсчёта таймаута до времени последнего ручного
// retry: после retry задача получает свежее окно и не падает повторно на
// ближайшем тике (см. state-reconciliation «Ручной повтор», MAJOR-1).
func (w *Worker) retriedFloor(d store.Download, basis time.Time) time.Time {
if r, ok := d.RetriedTime(); ok && r.After(basis) {
return r
}
return basis
}
// transition пишет новое состояние и логирует переход.
@@ -685,16 +735,22 @@ func (w *Worker) Retry(ctx context.Context, id string) error {
if d.State != store.StateFailed && d.State != store.StateStuck {
return fmt.Errorf("retry: download %s is %s, only failed/stuck are retriable", id, d.State)
}
// Если раздача уже жива в qBittorrent — перецепляемся к ней, повторный Add
// не нужен (и вреден: вслепую дублировал бы торрент). Add — только когда
// источника в qBittorrent нет. Базис таймаута берётся от added_on, поэтому
// возврат в downloading не роняет задачу снова на ближайшем тике.
alive := false
// Если раздача уже жива и ЗДОРОВА в qBittorrent — перецепляемся к ней,
// повторный Add не нужен (и вреден: вслепую дублировал бы торрент). Add —
// когда источника в qBittorrent нет ИЛИ он в состоянии ошибки: перецепка к
// сломанному торренту (error/missingFiles) бессмысленна — reconcile тут же
// вернул бы задачу в failed, поэтому пробуем повторно отдать источник
// (NIT-12). Базис таймаута сбрасывается ниже через retried_at, поэтому
// возврат в downloading не роняет задачу снова на ближайшем тике (MAJOR-1).
reAdd := true
if hashes := d.HashList(); len(hashes) > 0 {
_, alive, err = w.torrentByInfohash(ctx, hashes)
var t qbt.Torrent
var alive bool
t, alive, err = w.torrentByInfohash(ctx, hashes)
if err != nil {
return fmt.Errorf("retry: %w", err)
}
reAdd = !alive || classify(t.State) == classErrored
}
// Гард инварианта — ДО побочного эффекта в qBittorrent: пока задача лежала
// в failed, тем же infohash могла завладеть другая активная задача — тогда
@@ -705,7 +761,7 @@ func (w *Worker) Retry(ctx context.Context, id string) error {
}
return fmt.Errorf("retry: %w", err)
}
if !alive {
if reAdd {
// Добавляем заново по типу источника (magnet — ссылкой, torrent —
// сохранёнными байтами файлом). Rename при retry не выводим (namer здесь
// не зовём — имя уже могло быть выведено при первом добавлении).
@@ -729,6 +785,14 @@ func (w *Worker) Retry(ctx context.Context, id string) error {
return fmt.Errorf("retry: add to qbittorrent: %w", err)
}
}
// Сброс базиса отсчёта таймаутов (MAJOR-1): без него живой, но давно
// добавленный/простаивающий торрент снова упал бы по magnet_timeout/
// stuck_after на ближайшем тике. Best-effort: сбой лишь лишает свежего окна
// (WARN), сам retry уже состоялся.
if err := w.store.SetRetriedAt(ctx, id, w.now()); err != nil {
w.log.Warn("retry basis reset failed",
"capability", capReview, "download_id", id, "error", err)
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("download retried", "from", d.State)
return nil
}
+9
View File
@@ -211,6 +211,15 @@ func (f *fakeStore) SetSourceAddedAt(_ context.Context, id string, t time.Time)
return nil
}
func (f *fakeStore) SetRetriedAt(_ context.Context, id string, t time.Time) error {
d, ok := f.downloads[id]
if !ok {
return fmt.Errorf("download %s not found", id)
}
d.RetriedAt = store.NullString(store.FormatTime(t))
return nil
}
// --- Ф3-методы Store (заглушки; переопределяются в review_test.go) ---
func (f *fakeStore) CreateRecognition(_ context.Context, _ *store.Recognition, _ []string) (string, error) {