Приём: усыновление присутствующего в qBittorrent торрента вместо дубль-Add (409)

processCatched перед Add проверяет присутствие торрента в qBittorrent (один
листинг на тик): если раздача уже есть — усыновляем (promote catched→downloading
без повторного Add и без LLM-namer, имя из раздачи), иначе добавляем как раньше.
Это убирает бесконечный цикл дубль-Add → 409 → ретрай и лишние вызовы LLM.
Инвариант приёма «одна активная на infohash» делает различие «наш/чужой»
ненужным. source_type перечитывается под замком (сужение гонки апгрейда F6);
при недоступности qBittorrent тик пропускается без вызова LLM.

Дедуп на приёме (дубль на уже активную задачу) теперь отражается явным ответом
бота «дубль уже активной #id — добавление отменено».

Спека download-tracking обновлена (OpenSpec change заархивирован); закрыта
задача беклога review-f2-promote-without-add.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
av
2026-07-10 18:36:47 +03:00
co-authored by Claude Opus 4.8
parent 0c9421f4c1
commit b8657120fe
13 changed files with 544 additions and 47 deletions
+4 -3
View File
@@ -266,11 +266,12 @@ func (b *Bot) ingestAndReply(ctx context.Context, chatID int64, req ingest.Reque
b.send(chatID, opErr("Не удалось принять загрузку", res.DownloadID), nil)
return
}
msg := fmt.Sprintf("Принято #%s — добавляю в qBittorrent.", res.DownloadID)
if res.Deduplicated {
msg = fmt.Sprintf("Уже в работе #%s.", res.DownloadID)
// Дубль на уже активную задачу: новую загрузку не заводим, лишь сообщаем.
b.send(chatID, fmt.Sprintf("♻️ Дубль уже активной загрузки #%s — добавление отменено.", res.DownloadID), nil)
return
}
b.send(chatID, msg+"\nПозову, когда нужно подтверждение.", nil)
b.send(chatID, fmt.Sprintf("Принято #%s — добавляю в qBittorrent.\nПозову, когда нужно подтверждение.", res.DownloadID), nil)
}
const helpText = `jellybit-бот: пришлите magnet-ссылку, .torrent-файл или перешлите сообщение торрент-бота — поставлю на закачку.
+17
View File
@@ -162,6 +162,23 @@ func TestBot_IngestFromMagnet(t *testing.T) {
}
}
// Дедуп на приёме (дубль на уже активную задачу) → явный ответ «дубль …
// добавление отменено», а не «Принято».
func TestBot_IngestDeduplicated(t *testing.T) {
b, api, ing, _ := newTestBot(t, []int64{7})
ing.res = ingest.Result{DownloadID: tid, State: store.StateDownloading, Deduplicated: true}
b.handleMessage(context.Background(), msgFrom(7, "magnet:?xt=urn:btih:ABC"))
if len(api.sent) != 1 {
t.Fatalf("sent = %+v", api.sent)
}
txt := api.sent[0].text
if !strings.Contains(txt, "Дубль") || !strings.Contains(txt, tid) || strings.Contains(txt, "Принято") {
t.Errorf("ожидалось сообщение о дубле с #%s, got %q", tid, txt)
}
}
func TestBot_DeniesUnknownUser(t *testing.T) {
b, api, ing, _ := newTestBot(t, []int64{7})
b.handleMessage(context.Background(), msgFrom(999, "magnet:?xt=urn:btih:ABC"))
+81
View File
@@ -6,6 +6,7 @@ import (
"testing"
"time"
"git.vakhrushev.me/av/jellybit/internal/qbt"
"git.vakhrushev.me/av/jellybit/internal/store"
)
@@ -14,10 +15,12 @@ import (
type fakeNamer struct {
name string
gotContext string
calls int
onCall func()
}
func (f *fakeNamer) DeriveName(_ context.Context, contextText, _ string) string {
f.calls++
f.gotContext = contextText
if f.onCall != nil {
f.onCall()
@@ -77,6 +80,84 @@ func TestProcessCatchedAddsToQbit(t *testing.T) {
}
}
// Торрент пойманной загрузки уже присутствует в 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")
+78 -7
View File
@@ -387,10 +387,34 @@ func (w *Worker) processCatched(ctx context.Context) {
w.log.Warn("list catched failed", "capability", capIngest, "error", err)
return
}
if len(catched) == 0 {
return
}
// Снимок присутствия раздач в qBittorrent (один листинг на тик): по нему ДО
// вызова namer решаем, добавлять ли задачу вообще. Провал листинга —
// qBittorrent недоступен: пойманные в этот тик не трогаем (ни namer, ни Add),
// повтор на следующем; устойчивая недоступность отсекается предохранителем
// catch_timeout.
torrents, err := w.qbt.Torrents(ctx, "")
if err != nil {
w.log.Warn("list torrents for catched failed", "capability", capIngest, "error", err)
return
}
byHash := torrentsByHash(torrents)
for _, d := range catched {
cctx := w.scoped(ctx, capIngest, d.ID, d.PrimaryInfohash())
// Предохранитель: устойчивая невозможность добавить в qBittorrent.
// Торрент уже в qBittorrent — повторный Add не нужен (и вреден: qBittorrent
// отверг бы дубль, 409) и LLM-namer не зовём: усыновляем раздачу, доводя
// задачу до downloading. См. promoteExisting.
if t, ok := torrentFor(d, byHash); ok {
w.promoteExisting(cctx, d, t)
continue
}
// Предохранитель: устойчивая невозможность добавить в qBittorrent (раздачи
// в снимке нет и висит дольше catch_timeout).
if w.cfg.CatchTimeout > 0 {
if age, ok := w.catchedAge(d); ok && age > w.cfg.CatchTimeout {
w.mu.Lock()
@@ -406,9 +430,20 @@ func (w *Worker) processCatched(ctx context.Context) {
}
}
// Вне w.mu: сбор параметров добавления по типу источника (для torrent —
// чтение байтов), вывод имени (потенциально медленный LLM) и добавление.
hint, addReq, prepErr := w.sourceAddParts(cctx, d)
// Раздачи в qBittorrent нет — обычный путь добавления. Перечитываем запись
// под замком: (а) актуальный source_type (апгрейд magnet→torrent мог
// случиться после снятия списка catched — иначе добавили бы magnet из
// устаревшего снимка), (б) ре-валидация state=catched. Тяжёлые вызовы
// (чтение байтов, namer, Add) — вне замка.
w.mu.Lock()
cur, gerr := w.store.GetDownload(cctx, d.ID)
fresh := gerr == nil && cur != nil && cur.State == store.StateCatched
w.mu.Unlock()
if !fresh {
continue // отменили/пропала, пока шёл листинг — не трогаем
}
hint, addReq, prepErr := w.sourceAddParts(cctx, *cur)
if prepErr != nil {
// Байты torrent недоступны (не должно быть при штатном приёме) —
// остаёмся в catched, повтор на следующем тике.
@@ -417,13 +452,13 @@ func (w *Worker) processCatched(ctx context.Context) {
}
var rename string
if w.namer != nil {
rename = w.namer.DeriveName(cctx, d.Context, hint)
rename = w.namer.DeriveName(cctx, cur.Context, hint)
}
addReq.Rename = rename
addErr := w.qbt.Add(cctx, addReq)
if addErr != nil {
// Транзиентный сбой (qBit недоступен) — остаёмся в catched, повтор на
// следующем тике. Поведение вызова qBit уже залогировал клиент (ext.*).
// Транзиентный сбой (qBit отверг/недоступен) — остаёмся в catched,
// повтор на следующем тике. Вызов qBit уже залогировал клиент (ext.*).
logctx.From(cctx).Warn("catched add to qbittorrent failed, will retry", "error", addErr)
continue
}
@@ -441,6 +476,42 @@ func (w *Worker) processCatched(ctx context.Context) {
}
}
// promoteExisting усыновляет пойманную загрузку, чей торрент уже присутствует в
// qBittorrent (снимок тика): переводит catched → downloading БЕЗ повторного Add
// (иначе qBittorrent отверг бы дубль — 409 — и задача зациклилась бы) и без LLM.
// Имя берём из раздачи снимка (t.Name); у свежего magnet без метаданных (metaDL)
// оно может быть пустым — распознавание дольёт имя позже. Инвариант приёма
// гарантирует, что сюда доходит лишь загрузка без другой активной задачи на тот
// же infohash, поэтому присутствие раздачи трактуем как «усыновить и разложить»,
// а не как конфликт. Короткий DB-переход под w.mu; атомарный гард PromoteCatched
// (WHERE state='catched') сам отсекает гонку отмены, случившуюся, пока шёл листинг
// вне замка, — отдельный re-read не нужен.
func (w *Worker) promoteExisting(ctx context.Context, d store.Download, t qbt.Torrent) {
w.mu.Lock()
defer w.mu.Unlock()
if err := w.store.PromoteCatched(ctx, d.ID, t.Name); err != nil {
logctx.From(ctx).Info("catched promote skipped", "reason", err.Error())
return
}
logctx.From(ctx).Info("state transition", "from", store.StateCatched,
"to", store.StateDownloading, "reason", "already present in qbittorrent")
}
// torrentsByHash индексирует раздачи по каждому из их хешей (lowercase), как это
// делает Poll для своего снимка. Использует processCatched, чтобы проверить, есть
// ли торрент пойманной загрузки уже в qBittorrent.
func torrentsByHash(torrents []qbt.Torrent) map[string]qbt.Torrent {
byHash := make(map[string]qbt.Torrent, len(torrents)*2)
for _, t := range torrents {
for _, h := range []string{t.Hash, t.InfohashV1, t.InfohashV2} {
if h != "" {
byHash[strings.ToLower(h)] = t
}
}
}
return byHash
}
// catchedAge — возраст пойманной загрузки от created_at (у catched раздачи в
// qBittorrent ещё нет, added_on недоступен). ok=false — created_at не разобрать.
func (w *Worker) catchedAge(d store.Download) (time.Duration, bool) {
+4
View File
@@ -269,6 +269,7 @@ func (f *fakeStore) SetCandidateChosen(_ context.Context, _, _ string) error { r
type fakeQbt struct {
torrents []qbt.Torrent
torrentsErr error
onTorrents func() // вклинивается в момент листинга (симуляция гонки между снимком и re-read)
added []qbt.AddRequest
addErr error
files []qbt.File
@@ -288,6 +289,9 @@ type renameCall struct {
// регрессии: раздача, усыновлённая по тегу, имеет чужую категорию и не должна
// теряться при поиске по infohash.
func (f *fakeQbt) Torrents(_ context.Context, category string) ([]qbt.Torrent, error) {
if f.onTorrents != nil {
f.onTorrents()
}
if f.torrentsErr != nil {
return nil, f.torrentsErr
}