Приём: дедуп по target_missing/orphaned + стоп-кран «Закрыть»

Два дубля-близнеца на один инфохэш рождались, когда повторный приём
попадал на запись в target_missing: дедуп искал только активную задачу,
а target_missing терминален → заводилась новая загрузка, воркер усыновлял
уже присутствующий торрент и раскладывал его.

- Приём: критерий дедупа расширен до «блокирующей повторный приём» =
  активные ∪ {target_missing, orphaned}. Повторный приём такого инфохэша
  привязывается к существующей записи (спящей, без обращения к qBittorrent),
  а не плодит близнеца. Прочие терминальные (done/cancelled/failed/reverted/
  deleted) повторный приём не блокируют — осознанная свежая попытка. Новый
  read-метод FindReingestBlockingByInfohash (приоритет активной над desync);
  общий active-гард не тронут.
- Команда «Закрыть» (Dismiss) — универсальный стоп-кран из любого состояния,
  кроме deleted → cancelled (error_code=user_dismiss). Только меняет статус:
  файлы (в т.ч. хардлинки done/orphaned) и раздачу qBittorrent не трогает,
  в отличие от «Удалить». Веб — danger-зона внизу страницы; Telegram —
  кнопка с подтверждением; из cancelled — идемпотентный no-op.
- Транспорты при дедупе на desync-запись сообщают адресно (target_missing —
  привязать заново/закрыть; orphaned — закрыть и добавить заново); веб при
  дедупе ведёт на страницу существующей записи.

Спеки: ingest (дедуп), state-reconciliation (стоп-кран); граф переходов
допополнен рёбрами <терминал>→cancelled. OpenSpec change
dedup-target-missing-and-dismiss заархивирован.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
av
2026-07-10 20:15:37 +03:00
co-authored by Claude Opus 4.8
parent b8657120fe
commit 1369a9cabe
26 changed files with 1157 additions and 68 deletions
+2
View File
@@ -26,6 +26,7 @@ type actionReviewer struct {
stubReviewer
undoErr error
deleteErr error
dismissErr error
relinkErr error
rerecognizeErr error
refineErr error
@@ -34,6 +35,7 @@ type actionReviewer struct {
func (a actionReviewer) Undo(context.Context, string) error { return a.undoErr }
func (a actionReviewer) Delete(context.Context, string) error { return a.deleteErr }
func (a actionReviewer) Dismiss(context.Context, string) error { return a.dismissErr }
func (a actionReviewer) Relink(context.Context, string) error { return a.relinkErr }
func (a actionReviewer) Rerecognize(context.Context, string) error { return a.rerecognizeErr }
func (a actionReviewer) Refine(_ context.Context, _ string, hint string) error {
+7
View File
@@ -59,6 +59,11 @@ type downloadDetailView struct {
Relinkable bool
Retriable bool
Deletable bool // полное удаление доступно (done/orphaned/target_missing)
// Dismissable — доступен стоп-кран «Закрыть» (перевод в cancelled без
// действий над файлами/раздачей). Показываем в danger-зоне для терминальных,
// кроме deleted (строго терминален) и cancelled (уже закрыта, no-op); у
// не-терминальных ту же роль играет обычная «Отменить» — не дублируем.
Dismissable bool
}
// detailTitle — заголовок страницы просмотра: имя раздачи (display_name) →
@@ -119,6 +124,8 @@ func (s *server) buildDownloadView(id string, rd *worker.ReviewData) downloadDet
Retriable: d.State == store.StateFailed || d.State == store.StateStuck,
Deletable: d.State == store.StateDone || d.State == store.StateOrphaned ||
d.State == store.StateTargetMissing,
Dismissable: d.State.IsTerminal() &&
d.State != store.StateDeleted && d.State != store.StateCancelled,
}
// Дата добавления рядом с шапкой (source_added_at → фолбэк created_at,
// как в порядке и карточках списка); неразбираемое время просто опускаем.
+9
View File
@@ -140,6 +140,7 @@ func NewRouter(d Deps) (http.Handler, error) {
r.Post("/ui/downloads/{id}/undo", s.handleUndo)
r.Post("/ui/downloads/{id}/relink", s.handleRelink)
r.Post("/ui/downloads/{id}/delete", s.handleDelete)
r.Post("/ui/downloads/{id}/dismiss", s.handleDismiss)
// REST API.
r.Route("/api", func(r chi.Router) {
@@ -432,6 +433,14 @@ func (s *server) handleUIAdd(w http.ResponseWriter, r *http.Request) {
redirectErr(w, r, userErr(r, err, res.DownloadID))
return
}
if res.Deduplicated {
// Приём привязался к существующей записи (активной или «спящей» desync —
// target_missing/orphaned): ведём пользователя на её страницу, а не на
// список. Так видно, что нового не завели, и доступны действия записи
// (привязать заново / danger-зона «Закрыть»).
http.Redirect(w, r, "/download/"+res.DownloadID, http.StatusSeeOther)
return
}
http.Redirect(w, r, "/", http.StatusSeeOther)
}
+66
View File
@@ -490,6 +490,7 @@ type fakeReviewer struct {
deferred []string
undone []string
deleted []string
dismissed []string
relinked []string
rerecognized []string
cleared []string
@@ -532,6 +533,10 @@ func (f *fakeReviewer) Delete(_ context.Context, id string) error {
f.deleted = append(f.deleted, id)
return nil
}
func (f *fakeReviewer) Dismiss(_ context.Context, id string) error {
f.dismissed = append(f.dismissed, id)
return nil
}
func (f *fakeReviewer) Relink(_ context.Context, id string) error {
f.relinked = append(f.relinked, id)
return nil
@@ -614,6 +619,28 @@ func noRedirectClient() *http.Client {
}}
}
// Веб-приём при дедупе (в т.ч. на «спящую» desync-запись) ведёт на страницу
// существующей записи, а не на список — пользователь видит, что нового не завели.
func TestUIAddDeduplicatedRedirectsToRecord(t *testing.T) {
ing := &fakeIngestor{res: ingest.Result{DownloadID: tid, State: store.StateTargetMissing, Deduplicated: true}}
srv := newServer(t, httpapi.Deps{Ingestor: ing, Commander: &fakeCommander{}, Reader: &fakeReader{}})
req, _ := http.NewRequest(http.MethodPost, srv.URL+"/ui/downloads",
strings.NewReader("source=magnet:?xt=urn:btih:abc"))
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
resp, err := noRedirectClient().Do(req)
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusSeeOther {
t.Fatalf("status = %d, want 303", resp.StatusCode)
}
if loc := resp.Header.Get("Location"); loc != "/download/"+tid {
t.Errorf("Location = %q, want /download/%s", loc, tid)
}
}
func TestReviewRenders(t *testing.T) {
rv := &fakeReviewer{data: seriesReviewData()}
srv := newServer(t, httpapi.Deps{Ingestor: &fakeIngestor{}, Commander: &fakeCommander{},
@@ -765,6 +792,45 @@ func TestRefreshNameHTMXSwapsMain(t *testing.T) {
}
}
// Стоп-кран «Закрыть» рендерится в danger-зоне для терминального состояния
// (target_missing) и постит на /dismiss.
func TestDownloadPageShowsDismissButtonOnTerminal(t *testing.T) {
rd := seriesReviewData()
rd.Download.State = store.StateTargetMissing
rv := &fakeReviewer{data: rd}
srv := newServer(t, httpapi.Deps{Ingestor: &fakeIngestor{}, Commander: &fakeCommander{},
Reader: &fakeReader{}, Reviewer: rv})
resp, err := http.Get(srv.URL + "/download/" + tid)
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
if !strings.Contains(string(body), "/dismiss") || !strings.Contains(string(body), "Закрыть загрузку") {
t.Error("кнопка «Закрыть» не показана в danger-зоне для target_missing")
}
}
// POST /ui/downloads/{id}/dismiss вызывает Reviewer.Dismiss.
func TestUIDismiss(t *testing.T) {
rv := &fakeReviewer{data: seriesReviewData()}
srv := newServer(t, httpapi.Deps{Ingestor: &fakeIngestor{}, Commander: &fakeCommander{},
Reader: &fakeReader{}, Reviewer: rv})
resp, err := http.Post(srv.URL+"/ui/downloads/"+tid+"/dismiss", "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(rv.dismissed) != 1 || rv.dismissed[0] != tid {
t.Errorf("Dismiss вызван неверно: %v", rv.dismissed)
}
}
func TestDownloadPageShowsRefreshNameButtonOnDone(t *testing.T) {
// Кнопка гейтится наличием распознавания, не состоянием ревью: на терминальном
// done (есть план) она всё равно доступна.
+1
View File
@@ -50,6 +50,7 @@ func (stubReviewer) IgnoreFile(context.Context, string, string) error
func (stubReviewer) Defer(context.Context, string) error { return nil }
func (stubReviewer) Undo(context.Context, string) error { return nil }
func (stubReviewer) Delete(context.Context, string) error { return nil }
func (stubReviewer) Dismiss(context.Context, string) error { return nil }
func (stubReviewer) Relink(context.Context, string) error { return nil }
func (stubReviewer) Rerecognize(context.Context, string) error { return nil }
func (stubReviewer) ChooseCandidate(context.Context, string, string) error { return nil }
+12
View File
@@ -22,6 +22,7 @@ type Reviewer interface {
Defer(ctx context.Context, id string) error
Undo(ctx context.Context, id string) error
Delete(ctx context.Context, id string) error
Dismiss(ctx context.Context, id string) error
Relink(ctx context.Context, id string) error
Rerecognize(ctx context.Context, id string) error
ChooseCandidate(ctx context.Context, id, candidateID string) error
@@ -370,6 +371,17 @@ func (s *server) handleDelete(w http.ResponseWriter, r *http.Request) {
s.surfaceAction(w, r, id, s.deps.Reviewer.Delete(r.Context(), id))
}
// handleDismiss — универсальный стоп-кран: перевод задачи в cancelled только
// сменой статуса (файлы/раздачу не трогает), из danger-секции с подтверждением.
func (s *server) handleDismiss(w http.ResponseWriter, r *http.Request) {
id, err := pathID(r)
if err != nil {
redirectErr(w, r, "некорректный id")
return
}
s.surfaceAction(w, r, id, s.deps.Reviewer.Dismiss(r.Context(), id))
}
// handleRelink повторно привязывает откатанную задачу: перезапускает
// распознавание, задача пройдёт recognizing → review для подтверждения.
func (s *server) handleRelink(w http.ResponseWriter, r *http.Request) {
+22 -15
View File
@@ -23,9 +23,11 @@ const capIngest = "ingest"
// Store — нужная ingest часть хранилища.
type Store interface {
// FindActiveByInfohash — быстрый читающий дедуп-чек; авторитетная проверка —
// внутри CreateDownloadIfNoActive.
FindActiveByInfohash(ctx context.Context, hashes ...string) (*store.Download, error)
// FindReingestBlockingByInfohash — быстрый читающий дедуп-чек: активная задача
// ЛИБО удерживающая источник desync-запись (target_missing/orphaned). Активный
// инвариант «≤1 активной» авторитетно держит CreateDownloadIfNoActive; desync —
// устойчивый пред-рид, коротко замыкающий приём на возврат существующей записи.
FindReingestBlockingByInfohash(ctx context.Context, hashes ...string) (*store.Download, error)
// CreateDownloadIfNoActive атомарно проверяет инвариант «одна активная
// загрузка на infohash» и заводит задачу; вернувшаяся existing ≠ nil —
// дедуп на активную задачу (недостающие хеши вызова метод доносит сам).
@@ -87,15 +89,18 @@ func (s *Service) Ingest(ctx context.Context, req Request) (Result, error) {
log := s.log.With("capability", capIngest, "infohash", src.infohashes[0])
ctx = logctx.With(ctx, log)
// Быстрый дедуп-чек; авторитетная (атомарная) проверка — внутри
// CreateDownloadIfNoActive ниже. Дедуп — по ЛЮБОМУ из хешей источника:
// гибридный несёт и v1, и v2.
if existing, err := s.store.FindActiveByInfohash(ctx, src.infohashes...); err != nil {
// Быстрый дедуп-чек по ЛЮБОМУ из хешей источника (гибридный несёт и v1, и v2):
// активная задача ЛИБО удерживающая источник desync-запись
// (target_missing/orphaned) блокируют повторный приём. Для активной
// авторитетная (атомарная) проверка — внутри CreateDownloadIfNoActive ниже;
// desync-ветка сюда и завершается (в active-гард desync не заводим, чтобы не
// размыть инвариант «≤1 активной»).
if existing, err := s.store.FindReingestBlockingByInfohash(ctx, src.infohashes...); err != nil {
// Инфраструктурный сбой (БД) — операция приёма не выполнена: ERROR.
log.Error("ingest failed", "stage", "lookup-active", "error", err)
return Result{}, fmt.Errorf("ingest: lookup active: %w", err)
log.Error("ingest failed", "stage", "lookup-blocking", "error", err)
return Result{}, fmt.Errorf("ingest: lookup blocking: %w", err)
} else if existing != nil {
log.Info("download attached to active", "download_id", existing.ID, "state", existing.State)
log.Info("download attached", "download_id", existing.ID, "state", existing.State)
return s.attached(ctx, src, existing), nil
}
@@ -204,11 +209,13 @@ func mergeContext(userText, synth string) string {
return strings.Join(parts, "\n")
}
// attached — итог дедупа на быстром чеке: присоединились к уже активной
// задаче и доносим ей недостающие хеши источника (гибридный magnet мог
// принести хеш, которого задача ещё не знает; guarded-путь через
// CreateDownloadIfNoActive сюда не доходит). Донос — best-effort: конфликт
// хеша с другой активной задачей логируется, приём не валится.
// attached — итог дедупа на быстром чеке: присоединились к блокирующей записи
// (активной ЛИБО удерживающей источник desync — target_missing/orphaned) и
// доносим ей недостающие хеши источника (гибридный magnet мог принести хеш,
// которого задача ещё не знает; guarded-путь через CreateDownloadIfNoActive сюда
// не доходит). Для desync-записи состояние не меняем (возвращаем «спящей» — relink
// или закрытие делает пользователь). Донос — best-effort: конфликт хеша с другой
// активной задачей логируется, приём не валится.
func (s *Service) attached(ctx context.Context, src parsedSource, existing *store.Download) Result {
if len(src.infohashes) > len(existing.Infohashes) {
if err := s.store.AddInfohashes(ctx, existing.ID, src.infohashes); err != nil {
+26 -1
View File
@@ -26,7 +26,7 @@ type fakeStore struct {
upgradeUp bool // что вернуть из UpgradeCatchedMagnetToTorrent
}
func (f *fakeStore) FindActiveByInfohash(_ context.Context, _ ...string) (*store.Download, error) {
func (f *fakeStore) FindReingestBlockingByInfohash(_ context.Context, _ ...string) (*store.Download, error) {
return f.active, nil
}
@@ -159,6 +159,31 @@ func TestIngestIdempotent(t *testing.T) {
}
}
// Повторный приём привязывается к удерживающей источник desync-записи
// (target_missing/orphaned) вместо создания близнеца: возвращается существующая
// «спящей» (её состояние не меняется, к qBittorrent не ходим), Deduplicated=true.
func TestIngestAttachesToDesyncRecord(t *testing.T) {
for _, s := range []store.State{store.StateTargetMissing, store.StateOrphaned} {
t.Run(string(s), func(t *testing.T) {
existing := &store.Download{ID: "01hzzzexisting000000000000", State: s}
fs := &fakeStore{active: existing}
res, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet})
if err != nil {
t.Fatalf("Ingest: %v", err)
}
if !res.Deduplicated || res.DownloadID != existing.ID {
t.Errorf("ожидалось присоединение к desync-записи: %+v", res)
}
if res.State != s {
t.Errorf("состояние существующей записи должно вернуться как есть (%s), got %s", s, res.State)
}
if len(fs.created) != 0 {
t.Error("не должно создаваться новой задачи (близнеца)")
}
})
}
}
// Быстрый дедуп-путь доносит существующей задаче недостающие хеши
// гибридного magnet (иначе последующий приём по второму хешу создал бы
// вторую активную задачу).
+75 -8
View File
@@ -78,9 +78,14 @@ func (s State) IsTerminal() bool {
// (ActivateIfNoOtherActive): гейт графа ортогонален гарду терминальности в
// setState — граф говорит «ребро есть», гард «но не мимо ActivateIfNoOtherActive».
// Так, failed → downloading объявлено, но обычным SetDownloadState отклоняется.
// - cancelled/deferred — легальная цель из КАЖДОГО не-терминального состояния
// (Cancel/Defer проверяют лишь IsTerminal); инвариант закреплён тестом, а не
// ручной аккуратностью.
// - deferred — легальная цель из КАЖДОГО не-терминального состояния (Defer
// проверяет лишь IsTerminal); инвариант закреплён тестом, а не ручной
// аккуратностью.
// - cancelled — легальная цель из ЛЮБОГО состояния, кроме deleted: помимо
// Cancel из не-терминальных её даёт универсальный стоп-кран Dismiss, доступный
// и из терминальных (done/failed/reverted/target_missing/orphaned) — только
// смена статуса, файлы/раздачу не трогает (см. state-reconciliation «Ручное
// закрытие загрузки»). deleted строго терминален и цель cancelled не получает.
//
// Правка воркера, вводящая новое ребро, ОБЯЗАНА отразить его здесь — иначе
// setState отклонит переход (0 строк UPDATE → ошибка).
@@ -91,14 +96,14 @@ var allowedTransitions = map[State][]State{
StateRecognizing: {StateLinking, StateReview, StateCancelled, StateDeferred},
StateReview: {StateLinking, StateRecognizing, StateCancelled, StateDeferred, StateOrphaned, StateDeleted},
StateLinking: {StateDone, StateReview, StateFailed, StateCancelled, StateDeferred},
StateDone: {StateReverted, StateTargetMissing, StateOrphaned, StateDeleted},
StateDone: {StateReverted, StateTargetMissing, StateOrphaned, StateDeleted, StateCancelled},
StateDeferred: {StateLinking, StateRecognizing, StateCancelled, StateOrphaned, StateDeleted},
StateStuck: {StateDownloading, StateCompleted, StateCancelled, StateDeferred},
StateFailed: {StateDownloading, StateCompleted},
StateReverted: {StateRecognizing, StateOrphaned, StateDeleted},
StateFailed: {StateDownloading, StateCompleted, StateCancelled},
StateReverted: {StateRecognizing, StateOrphaned, StateDeleted, StateCancelled},
StateCancelled: {StateRecognizing, StateOrphaned, StateDeleted},
StateTargetMissing: {StateRecognizing, StateDone, StateOrphaned, StateDeleted},
StateOrphaned: {StateDone, StateTargetMissing, StateDeleted},
StateTargetMissing: {StateRecognizing, StateDone, StateOrphaned, StateDeleted, StateCancelled},
StateOrphaned: {StateDone, StateTargetMissing, StateDeleted, StateCancelled},
StateDeleted: nil, // окончательно терминально: сверка его не переоценивает
}
@@ -616,6 +621,68 @@ func (s *Store) FindActiveByInfohash(ctx context.Context, hashes ...string) (*Do
return d, nil
}
// reingestHoldingStates — desync-состояния, которые удерживают источник ради
// незакрытого намерения и потому БЛОКИРУЮТ повторный приём наравне с активными:
// target_missing (источник жив, ждёт relink) и orphaned (источник пропал, запись
// держит претензию на последнюю копию). Прочие терминальные (done/cancelled/
// failed/reverted/deleted) повторный приём НЕ блокируют — это осознанная свежая
// попытка. «Блокирующие» = активные (не-терминальные) reingestHoldingStates.
var reingestHoldingStates = []State{StateTargetMissing, StateOrphaned}
// FindReingestBlockingByInfohash возвращает задачу, блокирующую повторный приём
// любого из hashes: активную (строго не-терминальную) ЛИБО удерживающую источник
// desync-запись (target_missing/orphaned), приоритет — активной. Либо (nil, nil).
// Читающая основа расширенного дедупа приёма (пред-рид ДО создания); инвариант
// «≤1 активной на infohash» держит отдельный active-гард CreateDownloadIfNoActive,
// в который desync-состояния НЕ заводятся.
func (s *Store) FindReingestBlockingByInfohash(ctx context.Context, hashes ...string) (*Download, error) {
norm := normalizeHashes(hashes)
// Активная имеет приоритет: если по хешу есть и активная, и desync-запись
// (инвариант это допускает), присоединяемся к активной.
d, err := findActiveByInfohash(ctx, s.DB, norm, "")
if err != nil {
return nil, err
}
if d == nil {
d, err = findByInfohashInStates(ctx, s.DB, norm, reingestHoldingStates)
if err != nil {
return nil, err
}
}
if d != nil {
if err := attachInfohashesOne(ctx, s.DB, d); err != nil {
return nil, err
}
}
return d, nil
}
// findByInfohashInStates — выборка «задача по любому из хешей в одном из states»
// (позитивный фильтр `state IN (...)`, в отличие от findActiveByInfohash с
// `NOT IN terminalStates`). hashes уже нормализованы; хеши найденной загрузки НЕ
// подгружаются. Пустые hashes/states → (nil, nil).
func findByInfohashInStates(ctx context.Context, q sqlx.QueryerContext, hashes []string, states []State) (*Download, error) {
if len(hashes) == 0 || len(states) == 0 {
return nil, nil
}
var args []any
hashPh := placeholders(&args, hashes)
statePh := placeholders(&args, states)
query := `SELECT download.* FROM download
JOIN download_infohash dh ON dh.download_id = download.id
WHERE dh.infohash IN (` + hashPh + `) AND download.state IN (` + statePh + `)
ORDER BY download.id DESC LIMIT 1`
var d Download
err := sqlx.GetContext(ctx, q, &d, query, args...)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("find by infohash in states: %w", err)
}
return &d, nil
}
// findActiveByInfohash — общая выборка «активная задача по любому из хешей»
// (для guarded-методов — внутри их транзакции). hashes уже нормализованы;
// excludeID исключает саму проверяемую задачу (она может быть активной,
+65
View File
@@ -260,6 +260,71 @@ func TestFindActiveByInfohash_DesyncStatesNotActive(t *testing.T) {
}
}
// Дедуп повторного приёма блокируют не только активные, но и удерживающие
// источник desync-записи (target_missing/orphaned): по ним приём привязывается к
// существующей, а не плодит близнеца. Прочие терминальные (done/cancelled/failed/
// reverted/deleted) повторный приём НЕ блокируют — осознанная свежая попытка.
func TestFindReingestBlockingByInfohash(t *testing.T) {
ctx := context.Background()
blocking := []State{
StateCatched, StateDownloading, StateReview, // активные (примеры)
StateTargetMissing, StateOrphaned, // desync, удерживающие источник
}
for _, s := range blocking {
t.Run("blocking/"+string(s), func(t *testing.T) {
st := newTestStore(t)
ih := hashN(1)
id := mustCreate(t, st, ih)
forceState(t, st, id, s)
d, err := st.FindReingestBlockingByInfohash(ctx, ih)
if err != nil {
t.Fatal(err)
}
if d == nil || d.ID != id {
t.Fatalf("%s должна блокировать приём, получили %v", s, d)
}
if len(d.Infohashes) != 1 {
t.Fatalf("хеши не подгружены: %v", d.Infohashes)
}
})
}
nonBlocking := []State{
StateDone, StateCancelled, StateFailed, StateReverted, StateDeleted,
}
for _, s := range nonBlocking {
t.Run("non-blocking/"+string(s), func(t *testing.T) {
st := newTestStore(t)
ih := hashN(2)
id := mustCreate(t, st, ih)
forceState(t, st, id, s)
d, err := st.FindReingestBlockingByInfohash(ctx, ih)
if err != nil || d != nil {
t.Fatalf("%s блокировать приём не должна, получили (%v,%v)", s, d, err)
}
})
}
t.Run("priority-active-over-desync", func(t *testing.T) {
st := newTestStore(t)
ih := hashN(3)
// Старая запись ушла в target_missing (terminal освобождает хеш), затем по
// тому же хешу завелась новая активная — инвариант «≤1 активной» это
// допускает. Дедуп обязан присоединиться к активной, а не к desync.
oldID := mustCreate(t, st, ih)
forceState(t, st, oldID, StateTargetMissing)
newID := mustCreate(t, st, ih)
d, err := st.FindReingestBlockingByInfohash(ctx, ih)
if err != nil {
t.Fatal(err)
}
if d == nil || d.ID != newID {
t.Fatalf("ожидалась активная %s (приоритет над desync %s), получили %v", newID, oldID, d)
}
})
}
// Терминальное состояние освобождает infohash: тот же хеш заводится заново
// новой задачей (повторная закачка спустя время) — активность выводится
// только из state.
+29 -12
View File
@@ -74,19 +74,12 @@ func TestTransitionGraphWellFormed(t *testing.T) {
}
}
// Инвариант generic-команд Cancel/Defer: cancelled и deferred — легальная цель
// из КАЖДОГО не-терминального состояния (кроме самого deferred для deferred —
// это самопереход). Ловит класс дыры «забыли состояние» (напр. linking после
// краха процесса).
func TestCancelDeferReachableFromEveryNonTerminal(t *testing.T) {
// Инвариант Defer: deferred — легальная цель из КАЖДОГО не-терминального
// состояния (кроме самого deferred — это самопереход). Ловит класс дыры «забыли
// состояние» (напр. linking после краха процесса).
func TestDeferReachableFromEveryNonTerminal(t *testing.T) {
for _, s := range allStates {
if s.IsTerminal() {
continue
}
if !slices.Contains(transitionSources[StateCancelled], s) {
t.Errorf("%s → cancelled не легально (Cancel допускает любое не-терминальное)", s)
}
if s == StateDeferred {
if s.IsTerminal() || s == StateDeferred {
continue // deferred → deferred покрыт самопереходом
}
if !slices.Contains(transitionSources[StateDeferred], s) {
@@ -95,6 +88,22 @@ func TestCancelDeferReachableFromEveryNonTerminal(t *testing.T) {
}
}
// Инвариант универсального стоп-крана Dismiss: cancelled — легальная цель из
// ЛЮБОГО состояния, кроме deleted (строго терминален) и самого cancelled
// (самопереход). Шире инварианта Defer: покрывает и терминальные
// done/failed/reverted/target_missing/orphaned. НЕ объединять с проверкой
// deferred — у них разные множества источников.
func TestCancelledReachableFromEveryStateButDeleted(t *testing.T) {
for _, s := range allStates {
if s == StateDeleted || s == StateCancelled {
continue // deleted строго терминален; cancelled → cancelled — самопереход
}
if !slices.Contains(transitionSources[StateCancelled], s) {
t.Errorf("%s → cancelled не легально (Dismiss/Cancel допускают любое состояние, кроме deleted)", s)
}
}
}
// Объявленные не-revive рёбра проходят через SetDownloadState.
func TestSetStateAllowsDeclaredEdges(t *testing.T) {
edges := []struct{ from, to State }{
@@ -110,6 +119,14 @@ func TestSetStateAllowsDeclaredEdges(t *testing.T) {
{StateDone, StateReverted},
{StateStuck, StateCancelled},
{StateReview, StateDeferred},
// Стоп-кран Dismiss: терминал → cancelled идёт обычным SetDownloadState
// (цель cancelled терминальна → гард терминальности не мешает, revive не
// нужен).
{StateDone, StateCancelled},
{StateFailed, StateCancelled},
{StateReverted, StateCancelled},
{StateTargetMissing, StateCancelled},
{StateOrphaned, StateCancelled},
}
for i, e := range edges {
st := newTestStore(t)
+29 -2
View File
@@ -17,6 +17,7 @@ import (
"git.vakhrushev.me/av/jellybit/internal/ingest"
"git.vakhrushev.me/av/jellybit/internal/layout"
"git.vakhrushev.me/av/jellybit/internal/logging"
"git.vakhrushev.me/av/jellybit/internal/store"
"git.vakhrushev.me/av/jellybit/internal/worker"
)
@@ -47,6 +48,7 @@ type Reviewer interface {
Cancel(ctx context.Context, id string) error
Retry(ctx context.Context, id string) error
Delete(ctx context.Context, id string) error
Dismiss(ctx context.Context, id string) error
}
// Config — параметры бота.
@@ -267,8 +269,20 @@ func (b *Bot) ingestAndReply(ctx context.Context, chatID int64, req ingest.Reque
return
}
if res.Deduplicated {
// Дубль на уже активную задачу: новую загрузку не заводим, лишь сообщаем.
b.send(chatID, fmt.Sprintf("♻️ Дубль уже активной загрузки #%s — добавление отменено.", res.DownloadID), nil)
// Дубль: новую загрузку не заводим. Различаем активную задачу и «спящую»
// desync-запись (target_missing/orphaned) — у последней действие вперёд
// не «ждите», а «привяжите заново или закройте».
switch res.State {
case store.StateTargetMissing:
// Источник жив, цель удалена — из target_missing доступна перепривязка.
b.send(chatID, fmt.Sprintf("♻️ Этот торрент уже есть как запись #%s без цели — привяжите заново или закройте её.", res.DownloadID), nil)
case store.StateOrphaned:
// Источник пропал: relink из orphaned нет, рабочий путь — закрыть и
// добавить заново (тогда приём заведёт свежую загрузку).
b.send(chatID, fmt.Sprintf("♻️ Этот торрент уже есть как осиротевшая запись #%s — закройте её, затем добавьте заново.", res.DownloadID), nil)
default:
b.send(chatID, fmt.Sprintf("♻️ Дубль уже активной загрузки #%s — добавление отменено.", res.DownloadID), nil)
}
return
}
b.send(chatID, fmt.Sprintf("Принято #%s — добавляю в qBittorrent.\nПозову, когда нужно подтверждение.", res.DownloadID), nil)
@@ -327,6 +341,19 @@ func (b *Bot) handleCallback(ctx context.Context, cq *tgbotapi.CallbackQuery) {
case "delete_confirm":
err = b.reviewer.Delete(ctx, id)
note = "Удаляю…"
case "dismiss":
// Первый шаг: подтверждение. Стоп-кран лишь меняет статус (файлы/раздачу
// не трогает), но убирает запись из активной — подтверждаем сознательно.
b.answer(cq.ID, "")
b.editMarkup(chatID, msgID, b.dismissConfirmKeyboard(id))
return
case "dismiss_cancel":
b.answer(cq.ID, "Отменено")
b.refreshCard(ctx, chatID, msgID, id)
return
case "dismiss_confirm":
err = b.reviewer.Dismiss(ctx, id)
note = "Закрываю…"
case "type":
err = b.reviewer.SetType(ctx, id, val)
note = "Меняю тип…"
+30 -8
View File
@@ -62,14 +62,15 @@ func (f *fakeIngestor) Ingest(_ context.Context, req ingest.Request) (ingest.Res
}
type fakeReviewer struct {
data *worker.ReviewData
applied []string
refined map[string]string
typed map[string]string
deferred []string
canceled []string
retried []string
deleted []string
data *worker.ReviewData
applied []string
refined map[string]string
typed map[string]string
deferred []string
canceled []string
retried []string
deleted []string
dismissed []string
}
func (f *fakeReviewer) ReviewData(context.Context, string) (*worker.ReviewData, error) {
@@ -109,6 +110,10 @@ func (f *fakeReviewer) Delete(_ context.Context, id string) error {
f.deleted = append(f.deleted, id)
return nil
}
func (f *fakeReviewer) Dismiss(_ context.Context, id string) error {
f.dismissed = append(f.dismissed, id)
return nil
}
// tid — валидный lowercase-ULID (callback-data валидируется как ULID).
const tid = "01arz3ndektsv4rrffq69g5fav"
@@ -179,6 +184,23 @@ func TestBot_IngestDeduplicated(t *testing.T) {
}
}
// Дедуп на «спящую» desync-запись (target_missing) → сообщение зовёт привязать
// заново/закрыть, а не «дубль активной».
func TestBot_IngestDeduplicatedDesync(t *testing.T) {
b, api, ing, _ := newTestBot(t, []int64{7})
ing.res = ingest.Result{DownloadID: tid, State: store.StateTargetMissing, 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) {
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"))
+19 -3
View File
@@ -160,6 +160,8 @@ func (b *Bot) renderFailed(rd *worker.ReviewData) (string, *tgbotapi.InlineKeybo
func (b *Bot) retryKeyboard(id string) *tgbotapi.InlineKeyboardMarkup {
row := []tgbotapi.InlineKeyboardButton{
tgbotapi.NewInlineKeyboardButtonData("🔄 Повторить", "retry:"+id),
// Стоп-кран: закрыть зависшую задачу, не трогая файлы/раздачу.
tgbotapi.NewInlineKeyboardButtonData("✖️ Закрыть", "dismiss:"+id),
}
if url := b.reviewURL(id); url != "" {
row = append(row, tgbotapi.NewInlineKeyboardButtonURL("🌐 В вебе", url))
@@ -169,14 +171,18 @@ func (b *Bot) retryKeyboard(id string) *tgbotapi.InlineKeyboardMarkup {
}
// deletableKeyboard — клавиатура состояний, откуда доступно полное удаление
// (done/orphaned/target_missing): ссылка в веб (опц.) + «Удалить». Само удаление
// двухшаговое — кнопка ведёт на подтверждение (deleteConfirmKeyboard).
// (done/orphaned/target_missing): ссылка в веб (опц.) + «Закрыть» (стоп-кран, лишь
// статус) + «Удалить» (снос раздачи+файлов). Обе команды двухшаговые — кнопка
// ведёт на подтверждение.
func (b *Bot) deletableKeyboard(id string) *tgbotapi.InlineKeyboardMarkup {
var row []tgbotapi.InlineKeyboardButton
if url := b.reviewURL(id); url != "" {
row = append(row, tgbotapi.NewInlineKeyboardButtonURL("🌐 В вебе", url))
}
row = append(row, tgbotapi.NewInlineKeyboardButtonData("🗑 Удалить", "delete:"+id))
row = append(row,
tgbotapi.NewInlineKeyboardButtonData("✖️ Закрыть", "dismiss:"+id),
tgbotapi.NewInlineKeyboardButtonData("🗑 Удалить", "delete:"+id),
)
kb := tgbotapi.NewInlineKeyboardMarkup(tgbotapi.NewInlineKeyboardRow(row...))
return &kb
}
@@ -191,6 +197,16 @@ func (b *Bot) deleteConfirmKeyboard(id string) *tgbotapi.InlineKeyboardMarkup {
return &kb
}
// dismissConfirmKeyboard — шаг подтверждения закрытия (стоп-кран): перевод в
// «отменено» без действий над файлами/раздачей. Явное «Да» отделено от отмены.
func (b *Bot) dismissConfirmKeyboard(id string) *tgbotapi.InlineKeyboardMarkup {
kb := tgbotapi.NewInlineKeyboardMarkup(tgbotapi.NewInlineKeyboardRow(
tgbotapi.NewInlineKeyboardButtonData("✖️ Да, закрыть", "dismiss_confirm:"+id),
tgbotapi.NewInlineKeyboardButtonData("Отмена", "dismiss_cancel:"+id),
))
return &kb
}
func (b *Bot) webOnly(id string) *tgbotapi.InlineKeyboardMarkup {
url := b.reviewURL(id)
if url == "" {
+31
View File
@@ -897,6 +897,37 @@ func (w *Worker) Cancel(ctx context.Context, id string) (err error) {
return nil
}
// Dismiss — универсальный стоп-кран: переводит задачу в терминальный cancelled из
// ЛЮБОГО состояния, кроме deleted, ТОЛЬКО меняя статус. В отличие от Delete не
// трогает ни файлы (библиотечные хардлинки done/orphaned остаются на месте), ни
// раздачу в qBittorrent, ни цель; source-preflight не делает. Служит закрытием
// зависшей/спорной/лишней записи (в т.ч. дубля-близнеца в target_missing). Из
// cancelled — идемпотентный no-op БЕЗ setState: иначе перезаписал бы error_code,
// подменив причину прежнего Cancel/Dismiss. Помечает переход user_dismiss
// (отличает от reconcile и от штатного Cancel с пустым кодом).
func (w *Worker) Dismiss(ctx context.Context, id string) (err error) {
defer func() { w.logCmd(ctx, "dismiss", id, err) }()
w.mu.Lock()
defer w.mu.Unlock()
d, err := w.store.GetDownload(ctx, id)
if err != nil {
return fmt.Errorf("dismiss: %w", err)
}
if d.State == store.StateDeleted {
return fmt.Errorf("dismiss: download %s is deleted (strictly terminal): %w", id, ErrConflict)
}
if d.State == store.StateCancelled {
return nil // уже закрыта — no-op, error_code прежней отмены не трогаем
}
if err := w.store.SetDownloadState(ctx, id, store.StateCancelled, "user_dismiss", "закрыто пользователем"); err != nil {
return fmt.Errorf("dismiss: %w", err)
}
logctx.From(w.scoped(ctx, capReview, id, d.PrimaryInfohash())).Info("state transition",
"from", d.State, "to", store.StateCancelled, "code", "user_dismiss")
return nil
}
// Retry повторяет застрявшую/упавшую задачу: заново отдаёт источник в
// qBittorrent и возвращает в downloading.
func (w *Worker) Retry(ctx context.Context, id string) (err error) {
+52
View File
@@ -446,6 +446,58 @@ func TestCancel(t *testing.T) {
}
}
// Dismiss — стоп-кран: переводит в cancelled только сменой статуса, не трогая
// раздачу и файлы, из любого состояния, кроме deleted; из cancelled — no-op без
// перезаписи error_code.
func TestDismiss(t *testing.T) {
t.Run("from-done-keeps-source-and-files", func(t *testing.T) {
st := oneDownloading("541adcff3b6dd5dba7088ea83317d9d6fac331d6", timeRecent)
st.downloads["1"].State = store.StateDone
qb := &fakeQbt{}
w := newTestWorker(st, qb)
if err := w.Dismiss(context.Background(), "1"); err != nil {
t.Fatalf("Dismiss: %v", err)
}
if st.downloads["1"].State != store.StateCancelled {
t.Errorf("state = %q, want cancelled", st.downloads["1"].State)
}
if got := st.downloads["1"].ErrorCode.String; got != "user_dismiss" {
t.Errorf("error_code = %q, want user_dismiss", got)
}
if len(qb.deleted) != 0 {
t.Errorf("раздачу трогать не должны, Delete вызван %d раз", len(qb.deleted))
}
})
t.Run("from-deleted-rejected", func(t *testing.T) {
st := oneDownloading("541adcff3b6dd5dba7088ea83317d9d6fac331d6", timeRecent)
st.downloads["1"].State = store.StateDeleted
w := newTestWorker(st, &fakeQbt{})
if err := w.Dismiss(context.Background(), "1"); err == nil {
t.Error("ожидалась ошибка Dismiss из deleted")
}
if st.downloads["1"].State != store.StateDeleted {
t.Errorf("state = %q, want deleted (не изменилось)", st.downloads["1"].State)
}
})
t.Run("from-cancelled-noop-keeps-code", func(t *testing.T) {
st := oneDownloading("541adcff3b6dd5dba7088ea83317d9d6fac331d6", timeRecent)
st.downloads["1"].State = store.StateCancelled
st.downloads["1"].ErrorCode = store.NullString("prior_reason")
w := newTestWorker(st, &fakeQbt{})
if err := w.Dismiss(context.Background(), "1"); err != nil {
t.Fatalf("Dismiss no-op: %v", err)
}
if got := st.downloads["1"].ErrorCode.String; got != "prior_reason" {
t.Errorf("error_code = %q, no-op не должен его переписывать", got)
}
if len(st.transitions) != 0 {
t.Errorf("no-op не должен писать переход, got %d", len(st.transitions))
}
})
}
func TestRetry(t *testing.T) {
st := oneDownloading("541adcff3b6dd5dba7088ea83317d9d6fac331d6", timeRecent)
st.downloads["1"].State = store.StateStuck