diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index 394da52..0ecad06 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -35,6 +35,11 @@ type Store interface { // AddInfohashes доносит задаче недостающие хеши (guarded). Нужен на // быстром дедуп-пути, который не доходит до CreateDownloadIfNoActive. AddInfohashes(ctx context.Context, downloadID string, hashes []string) error + // UpgradeCatchedMagnetToTorrent при дедупе входящих байтов `.torrent` на + // пойманную (`catched`) magnet-задачу сохраняет байты и меняет source_type + // на torrent (атомарно). No-op, если задача уже добавлена/не magnet. Чинит + // magnet закрытого трекера, который иначе застрянет в metaDL. + UpgradeCatchedMagnetToTorrent(ctx context.Context, downloadID string, torrentBlob []byte) (bool, error) } // Service — реализация быстрого приёма. @@ -109,14 +114,9 @@ func (s *Service) Ingest(ctx context.Context, req Request) (Result, error) { if existing != nil { // Гонка с параллельным приёмом/discover: активная задача появилась // после быстрого чека — присоединяемся к ней (хеши донёс сам - // CreateDownloadIfNoActive). + // CreateDownloadIfNoActive; existing уже с подгруженными хешами). log.Info("download attached to active", "download_id", existing.ID, "state", existing.State) - return Result{ - DownloadID: existing.ID, - Infohashes: src.infohashes, - State: existing.State, - Deduplicated: true, - }, nil + return s.attached(ctx, src, existing), nil } log.Info("download catched", "download_id", d.ID) @@ -208,6 +208,17 @@ func (s *Service) attached(ctx context.Context, src parsedSource, existing *stor logctx.FromOr(ctx, s.log).Warn("ingest top-up infohashes failed", "error", err) } } + // F6: входящее — байты `.torrent`, а активная задача поймана как magnet и + // ещё не отдана в qBittorrent (catched) → сохраняем байты и переключаем + // источник на torrent, чтобы воркер добавил раздачу файлом (magnet + // закрытого трекера иначе застрянет в metaDL). Best-effort: не валит приём. + if src.sourceType == store.SourceTorrent && len(src.torrentBlob) > 0 { + if upgraded, err := s.store.UpgradeCatchedMagnetToTorrent(ctx, existing.ID, src.torrentBlob); err != nil { + logctx.FromOr(ctx, s.log).Warn("ingest torrent upgrade failed", "error", err) + } else if upgraded { + logctx.FromOr(ctx, s.log).Info("catched magnet upgraded to torrent", "download_id", existing.ID) + } + } return Result{ DownloadID: existing.ID, Infohashes: src.infohashes, diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 74f1b1a..607c146 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -16,11 +16,14 @@ const sampleMagnet = "magnet:?xt=urn:btih:541ADCFF3B6DD5DBA7088EA83317D9D6FAC331 const sampleInfohash = "541adcff3b6dd5dba7088ea83317d9d6fac331d6" type fakeStore struct { - active *store.Download - created []store.Download - hashes [][]string - blobs [][]byte - toppedUp []string + active *store.Download + created []store.Download + hashes [][]string + blobs [][]byte + toppedUp []string + upgradeID string // downloadID последнего вызова UpgradeCatchedMagnetToTorrent + upgradeBlob []byte // байты, переданные в апгрейд + upgradeUp bool // что вернуть из UpgradeCatchedMagnetToTorrent } func (f *fakeStore) FindActiveByInfohash(_ context.Context, _ ...string) (*store.Download, error) { @@ -44,6 +47,12 @@ func (f *fakeStore) AddInfohashes(_ context.Context, id string, hashes []string) return nil } +func (f *fakeStore) UpgradeCatchedMagnetToTorrent(_ context.Context, id string, blob []byte) (bool, error) { + f.upgradeID = id + f.upgradeBlob = blob + return f.upgradeUp, nil +} + func newService(st Store) *Service { return New(st, slog.New(slog.NewTextHandler(io.Discard, nil))) } diff --git a/internal/ingest/torrent_ingest_test.go b/internal/ingest/torrent_ingest_test.go index 186c0a6..895bbd4 100644 --- a/internal/ingest/torrent_ingest_test.go +++ b/internal/ingest/torrent_ingest_test.go @@ -83,6 +83,45 @@ func TestIngestTorrentDedupNoBlob(t *testing.T) { } } +// F6: дедуп .torrent на пойманную (catched) magnet-задачу вызывает апгрейд — +// сохранение байтов и смену источника (magnet закрытого трекера иначе застрянет +// в metaDL). +func TestIngestTorrentUpgradesCatchedMagnet(t *testing.T) { + data, hash := buildTorrent(t, "Dune", "http://t/ann") + existing := &store.Download{ + ID: "cm", + State: store.StateCatched, + SourceType: store.SourceMagnet, + Infohashes: []store.Infohash{{Infohash: hash, Kind: store.HashV1}}, + } + fs := &fakeStore{active: existing, upgradeUp: true} + res, err := newService(fs).Ingest(context.Background(), Request{TorrentData: data}) + if err != nil { + t.Fatalf("Ingest: %v", err) + } + if !res.Deduplicated || res.DownloadID != "cm" { + t.Errorf("ожидался дедуп на cm, res = %+v", res) + } + if fs.upgradeID != "cm" { + t.Errorf("апгрейд не вызван для существующей задачи (upgradeID=%q)", fs.upgradeID) + } + if !bytes.Equal(fs.upgradeBlob, data) { + t.Errorf("в апгрейд переданы не те байты") + } +} + +// Дедуп magnet-ссылки (не .torrent) апгрейд не вызывает — нечего сохранять. +func TestIngestMagnetDedupNoUpgrade(t *testing.T) { + existing := &store.Download{ID: "cm", State: store.StateCatched, SourceType: store.SourceMagnet} + fs := &fakeStore{active: existing} + if _, err := newService(fs).Ingest(context.Background(), Request{Source: sampleMagnet}); err != nil { + t.Fatalf("Ingest: %v", err) + } + if fs.upgradeID != "" { + t.Errorf("апгрейд не должен вызываться для magnet-дедупа") + } +} + func TestIngestTorrentTooLarge(t *testing.T) { fs := &fakeStore{} big := make([]byte, MaxTorrentSize+1) diff --git a/internal/store/download.go b/internal/store/download.go index bb670da..174aa25 100644 --- a/internal/store/download.go +++ b/internal/store/download.go @@ -272,8 +272,18 @@ func (s *Store) CreateDownloadIfNoActive(ctx context.Context, d *Download, hashe // Дедуп нашёл активного владельца по одному из хешей — остальные хеши // norm принадлежат тому же торренту (гибридный magnet): дописываем // недостающие, иначе второй хеш молча теряется и последующий приём по - // нему создал бы вторую активную задачу. + // нему создал бы вторую активную задачу. Дозапись — под тем же пер-хеш + // гардом владения, что и AddInfohashes: хеш, которым владеет ДРУГАЯ + // активная задача (split-identity/крафт-магнет), не дописываем, иначе + // две активные владели бы одним инфохэшем (нарушение инварианта). for _, h := range norm { + owner, err := findActiveByInfohash(ctx, tx, []string{h}, existing.ID) + if err != nil { + return nil, fmt.Errorf("create download: %w", err) + } + if owner != nil { + continue // чужой активный владелец — не крадём хеш + } if _, err := tx.ExecContext(ctx, `INSERT OR IGNORE INTO download_infohash (download_id, infohash, kind, created_at) VALUES (?, ?, ?, ?)`, existing.ID, h, HashKind(h), now); err != nil { @@ -335,6 +345,57 @@ func (s *Store) GetTorrentData(ctx context.Context, downloadID string) ([]byte, return data, nil } +// UpgradeCatchedMagnetToTorrent апгрейдит пойманную magnet-задачу до torrent +// при дедупе входящих байтов `.torrent`: в одной write-транзакции сохраняет +// байты и меняет source_type magnet→torrent, но ТОЛЬКО пока задача в `catched` +// (воркер источник ещё не отдал в qBittorrent) и её source_type всё ещё +// `magnet`. Это целевое исключение из правила «при дедупе байты не сохраняем»: +// magnet на закрытом трекере без DHT метаданные не докачает и застрянет в +// metaDL, а поданный пользователем `.torrent` их несёт. Возвращает true, если +// апгрейд применён; (false, nil) — если задача не подходит (уже добавлена, +// отменена или не magnet) либо байты пусты (защитный no-op). Гард через +// RowsAffected UPDATE — та же ре-валидация состояния, что у PromoteCatched. +func (s *Store) UpgradeCatchedMagnetToTorrent(ctx context.Context, downloadID string, torrentBlob []byte) (bool, error) { + if len(torrentBlob) == 0 { + return false, nil + } + tx, err := s.DB.BeginTxx(ctx, nil) + if err != nil { + return false, fmt.Errorf("upgrade %s to torrent: begin tx: %w", downloadID, err) + } + defer func() { _ = tx.Rollback() }() + + res, err := tx.ExecContext(ctx, ` +UPDATE download +SET source_type = ?, updated_at = ? +WHERE id = ? AND source_type = ? AND state = ?`, + string(SourceTorrent), FormatTime(Now()), downloadID, string(SourceMagnet), string(StateCatched)) + if err != nil { + return false, fmt.Errorf("upgrade %s to torrent: %w", downloadID, err) + } + n, err := res.RowsAffected() + if err != nil { + return false, fmt.Errorf("upgrade %s to torrent: %w", downloadID, err) + } + if n == 0 { + // Не catched-magnet (уже добавлена/отменена/torrent) — апгрейд не нужен. + if err := tx.Commit(); err != nil { + return false, fmt.Errorf("upgrade %s to torrent: commit: %w", downloadID, err) + } + return false, nil + } + // У magnet-задачи блоба не было; OR REPLACE — страховка идемпотентности. + if _, err := tx.ExecContext(ctx, + `INSERT OR REPLACE INTO download_torrent (download_id, data) VALUES (?, ?)`, + downloadID, torrentBlob); err != nil { + return false, fmt.Errorf("upgrade %s to torrent: store blob: %w", downloadID, err) + } + if err := tx.Commit(); err != nil { + return false, fmt.Errorf("upgrade %s to torrent: commit: %w", downloadID, err) + } + return true, nil +} + // ActivateIfNoOtherActive атомарно возвращает загрузку в активное состояние // (retry/восстановление сверкой/relink): в одной write-транзакции проверяет, // что никакая ДРУГАЯ активная загрузка не владеет любым из хешей этой, и diff --git a/internal/store/download_test.go b/internal/store/download_test.go index bd587ac..b1d8688 100644 --- a/internal/store/download_test.go +++ b/internal/store/download_test.go @@ -3,6 +3,7 @@ package store import ( "context" "errors" + "slices" "strings" "sync" "testing" @@ -167,6 +168,50 @@ func TestFindActiveByInfohash(t *testing.T) { } } +// F1: дедуп-ветка CreateDownloadIfNoActive доносит недостающие хеши под пер-хеш +// гардом владения — хеш, которым владеет ДРУГАЯ активная задача (split-identity +// гибрида), красть не должна, иначе две активные владели бы одним инфохэшем. +func TestDedupTopUpDoesNotStealForeignHash(t *testing.T) { + st := newTestStore(t) + ctx := context.Background() + v1 := hashN(1) + v2 := hashN(2) + + // A владеет v1, B владеет v2 — две активные задачи одного гибрида (split). + idA := mustCreate(t, st, v1) + idB := mustCreate(t, st, v2) + + // Приём гибрида {v1, v2} → дедуп на одну из активных; чужой хеш не крадём. + existing, err := st.CreateDownloadIfNoActive(ctx, newDownloading(), []string{v1, v2}, nil) + if err != nil { + t.Fatalf("create: %v", err) + } + if existing == nil { + t.Fatal("ожидался дедуп на активную задачу") + } + + assertHashes(t, st, idA, v1) // A по-прежнему владеет только v1 + assertHashes(t, st, idB, v2) // B по-прежнему владеет только v2 +} + +// assertHashes проверяет точный набор инфохэшей загрузки. +func assertHashes(t *testing.T, st *Store, id string, want ...string) { + t.Helper() + d, err := st.GetDownload(context.Background(), id) + if err != nil { + t.Fatalf("get %s: %v", id, err) + } + got := d.HashList() + if len(got) != len(want) { + t.Fatalf("%s hashes = %v, want %v", id, got, want) + } + for _, w := range want { + if !slices.Contains(got, w) { + t.Fatalf("%s hashes = %v, want содержит %s", id, got, w) + } + } +} + // Состояния рассинхрона (target_missing/orphaned/deleted) терминальны: задача // в них не должна считаться «активной» (иначе relink/ingest-дедуп решат, что // для infohash уже есть активная задача). Регрессия: FindActiveByInfohash и diff --git a/internal/store/torrent_blob_test.go b/internal/store/torrent_blob_test.go index 79412ea..cf45d96 100644 --- a/internal/store/torrent_blob_test.go +++ b/internal/store/torrent_blob_test.go @@ -65,3 +65,110 @@ func TestTorrentBlobMissing(t *testing.T) { t.Errorf("ожидался ErrNotFound, got %v", err) } } + +// F6: апгрейд пойманного magnet до torrent сохраняет байты и меняет source_type. +func TestUpgradeCatchedMagnetToTorrent(t *testing.T) { + st := newTestStore(t) + ctx := context.Background() + d := &Download{SourceType: SourceMagnet, SourceRef: "magnet:?xt=urn:btih:" + hashN(1), State: StateCatched} + if _, err := st.CreateDownloadIfNoActive(ctx, d, []string{hashN(1)}, nil); err != nil { + t.Fatalf("create: %v", err) + } + + blob := []byte("d8:announce…real torrent") + upgraded, err := st.UpgradeCatchedMagnetToTorrent(ctx, d.ID, blob) + if err != nil || !upgraded { + t.Fatalf("upgrade: upgraded=%v err=%v", upgraded, err) + } + + got, err := st.GetDownload(ctx, d.ID) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.SourceType != SourceTorrent { + t.Errorf("source_type = %q, want torrent", got.SourceType) + } + data, err := st.GetTorrentData(ctx, d.ID) + if err != nil { + t.Fatalf("get torrent data: %v", err) + } + if string(data) != string(blob) { + t.Errorf("blob = %q, want %q", data, blob) + } +} + +// Апгрейд применим только в catched: уже добавленный (downloading) magnet не +// трогаем — его судьба решается retry/сверкой, а не приёмом. +func TestUpgradeSkippedWhenNotCatched(t *testing.T) { + st := newTestStore(t) + ctx := context.Background() + d := &Download{SourceType: SourceMagnet, SourceRef: "m", State: StateDownloading} + if _, err := st.CreateDownloadIfNoActive(ctx, d, []string{hashN(2)}, nil); err != nil { + t.Fatalf("create: %v", err) + } + + upgraded, err := st.UpgradeCatchedMagnetToTorrent(ctx, d.ID, []byte("bytes")) + if err != nil { + t.Fatalf("upgrade: %v", err) + } + if upgraded { + t.Error("downloading-magnet апгрейдить не должны") + } + got, err := st.GetDownload(ctx, d.ID) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.SourceType != SourceMagnet { + t.Errorf("source_type сменился на %q, а не должен", got.SourceType) + } + if _, err := st.GetTorrentData(ctx, d.ID); !errors.Is(err, ErrNotFound) { + t.Errorf("байты не должны сохраняться, got err=%v", err) + } +} + +// Не-magnet источник апгрейд не трогает (гард source_type=magnet), байты не +// перезаписывает. +func TestUpgradeSkippedForTorrentSource(t *testing.T) { + st := newTestStore(t) + ctx := context.Background() + d := &Download{SourceType: SourceTorrent, SourceRef: "Rel", State: StateCatched} + if _, err := st.CreateDownloadIfNoActive(ctx, d, []string{hashN(3)}, []byte("orig")); err != nil { + t.Fatalf("create: %v", err) + } + + upgraded, err := st.UpgradeCatchedMagnetToTorrent(ctx, d.ID, []byte("new")) + if err != nil { + t.Fatalf("upgrade: %v", err) + } + if upgraded { + t.Error("torrent-источник апгрейдить не нужно") + } + data, err := st.GetTorrentData(ctx, d.ID) + if err != nil { + t.Fatalf("get torrent data: %v", err) + } + if string(data) != "orig" { + t.Errorf("байты перезаписаны: %q", data) + } +} + +// Пустые байты — защитный no-op. +func TestUpgradeEmptyBlobNoOp(t *testing.T) { + st := newTestStore(t) + ctx := context.Background() + d := &Download{SourceType: SourceMagnet, SourceRef: "m", State: StateCatched} + if _, err := st.CreateDownloadIfNoActive(ctx, d, []string{hashN(4)}, nil); err != nil { + t.Fatalf("create: %v", err) + } + upgraded, err := st.UpgradeCatchedMagnetToTorrent(ctx, d.ID, nil) + if err != nil || upgraded { + t.Fatalf("пустой блоб: upgraded=%v err=%v", upgraded, err) + } + got, err := st.GetDownload(ctx, d.ID) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.SourceType != SourceMagnet { + t.Errorf("source_type сменился на %q при пустом блобе", got.SourceType) + } +} diff --git a/openspec/changes/ingest-dedup-integrity/.openspec.yaml b/openspec/changes/ingest-dedup-integrity/.openspec.yaml new file mode 100644 index 0000000..8cceb8d --- /dev/null +++ b/openspec/changes/ingest-dedup-integrity/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-07-08 diff --git a/openspec/changes/ingest-dedup-integrity/design.md b/openspec/changes/ingest-dedup-integrity/design.md new file mode 100644 index 0000000..90e12d7 --- /dev/null +++ b/openspec/changes/ingest-dedup-integrity/design.md @@ -0,0 +1,79 @@ +## Context + +Приём дедуплицирует входящий источник по активной задаче двумя путями: + +1. **Быстрый чек** — `Ingest` вызывает `FindActiveByInfohash`, при попадании + уходит в `attached()` (дозапись недостающих хешей через `AddInfohashes`). + Это основной путь. +2. **Гонка** — быстрый чек пуст, но пока выводили имя/готовили запись, + активная задача появилась; `CreateDownloadIfNoActive` внутри своей + транзакции находит её и возвращает как дедуп (дозапись хешей внутри той же + tx). + +F1 живёт в пути 2 (неохраняемая дозапись). F6 задевает ОБА пути: в реальном +сценарии первым отрабатывает быстрый чек (путь 1), поэтому апгрейд обязан +работать и там. + +## F1 — пер-хеш гард дозаписи в дедуп-ветке + +**Решение.** В дедуп-ветке `CreateDownloadIfNoActive` заменяем безусловный +цикл `INSERT OR IGNORE` на тот же гард, что в `AddInfohashes`: для каждого +хеша `h` проверяем `findActiveByInfohash(ctx, tx, {h}, existing.ID)`; если +другая активная задача владеет `h` — пропускаем (не дописываем), иначе +`INSERT OR IGNORE`. Транзакция уже открыта, `excludeID = existing.ID` +исключает саму дедуп-цель (её собственные хеши не конфликтуют с ней самой). + +**Почему пропуск, а не ошибка.** Дедуп по контракту не падает — возвращает +существующую задачу. Конфликтный хеш принадлежит другой активной задаче; +молча его не трогаем — это и есть сохранение инварианта. В отличие от +`AddInfohashes`, который сигналит `ErrInfohashTaken` вызывающему (там это +осмысленно), у дедуп-ветки наблюдаемого канала ошибки нет и он не нужен: +поведение — «присоединиться к найденной задаче, чужое не красть». + +## F6 — апгрейд catched-magnet до torrent + +**Решение.** Новый guarded-метод хранилища: + +``` +UpgradeCatchedMagnetToTorrent(ctx, downloadID string, torrentBlob []byte) (bool, error) +``` + +в одной write-транзакции: + +1. Пустой `torrentBlob` → `(false, nil)` (защитный no-op). +2. Гардированный UPDATE: + `UPDATE download SET source_type='torrent', updated_at=? + WHERE id=? AND source_type='magnet' AND state='catched'`. + `RowsAffected==0` → задача не подходит (уже `downloading`/отменена/не + magnet) → коммит без эффекта, `(false, nil)`. +3. `RowsAffected==1` → `INSERT OR REPLACE INTO download_torrent(download_id, + data)` (у magnet блоба нет; `OR REPLACE` — страховка идемпотентности), + коммит, `(true, nil)`. + +**Почему гард `state='catched'`.** Апгрейд имеет смысл только пока worker ещё +не отдал источник в qBittorrent. В `catched` worker на шаге добавления +выбирает способ по `source_type` (`internal/worker/worker.go:sourceAddParts`): +после смены на `torrent` он добавит файлом — метаданные приедут сразу. Если +задача уже `downloading`, magnet давно в qBittorrent (застрял в metaDL), и +смена `source_type` его не переотдаст; это отдельная забота retry/desync, не +приёма. Тот же паттерн ре-валидации, что у `PromoteCatched`: если задачу +успели отменить, UPDATE не заденет строк и апгрейд просто не применится. + +**Где вызываем.** Оба дедуп-пути `Ingest` сводим к `attached()` (в ветке +гонки `existing` от `CreateDownloadIfNoActive` уже с подгруженными хешами, так +что переиспользование безопасно). В `attached()` после дозаписи хешей: если +входящий источник — `torrent` и есть байты, зовём +`UpgradeCatchedMagnetToTorrent`. Вызов best-effort: ошибка/неуспех логируются +`Warn`/`Info`, приём не валится (как и дозапись хешей). Результат приёма +по-прежнему несёт `State` существующей задачи (остаётся `catched`) — меняется +лишь способ будущего добавления. + +**Идемпотентность и повтор.** Повторный `.torrent` того же хеша: первый +апгрейд перевёл задачу в `torrent`, гард `source_type='magnet'` на втором даст +`RowsAffected==0` → no-op. Байты уже сохранены при создании torrent-ветки — +`OR REPLACE` перезапишет теми же данными без вреда. + +## Границы + +Схему не трогаем: `download_torrent` и колонка `source_type` уже есть. ER-схема +`docs/specs/database.md` без изменений. Апгрейд — операция над данными. diff --git a/openspec/changes/ingest-dedup-integrity/proposal.md b/openspec/changes/ingest-dedup-integrity/proposal.md new file mode 100644 index 0000000..e9b771f --- /dev/null +++ b/openspec/changes/ingest-dedup-integrity/proposal.md @@ -0,0 +1,56 @@ +## Why + +Ревью приёма (Fable, 2026-07-08) вскрыло два дефекта в дедуп-ветках приёма, +оба про инвариант «≤1 активная загрузка на infohash» и про сохранность +источника: + +- **F1 — неохраняемая дозапись хешей.** Дедуп-ветка + `CreateDownloadIfNoActive` (`internal/store/download.go`) безусловно + дописывает ВСЕ хеши входящего источника в найденную активную задачу + (`INSERT OR IGNORE`) без пер-хеш гарда владения — в отличие от + `AddInfohashes`, где гард есть. Это единственная неохраняемая запись хешей, + и она в авторитетном методе инварианта. Сценарий: активная A владеет `v1`, + активная B владеет `v2` того же гибридного торрента; приём гибрида + `{v1,v2}`, дедупнувшись на B, допишет `v1` в B → две активные владеют `v1`. + Инвариант нарушен. + +- **F6 — потерянный upgrade-путь `.torrent` поверх magnet.** Пользователь + сначала ловит magnet с закрытого трекера (задача `catched`, + `source_type=magnet`), понимает, что без DHT метаданные не докачаются, и + грузит правильный `.torrent`. Приём дедупит по infohash на magnet-задачу; по + текущей спеке байты при дедупе НЕ сохраняются, `source_type` остаётся + `magnet`. Worker добавляет раздачу по magnet-URL → вечный metaDL → failed. + Ровно тот артефакт, который бы починил загрузку, выбрасывается с «уже в + работе». + +## What Changes + +- **F1:** дедуп-ветка `CreateDownloadIfNoActive` применяет тот же пер-хеш + гард, что и `AddInfohashes`: хеш, которым владеет ДРУГАЯ активная задача, + не дописывается. Транзакция уже открыта — правка внутри неё. + +- **F6:** при дедупе, где входящее — байты `.torrent`, а активная задача + поймана как `magnet` и ещё не добавлена в qBittorrent (состояние + `catched`), система сохраняет байты и меняет `source_type` на `torrent` в + одной транзакции. Тогда worker добавит раздачу файлом и метаданные не + придётся докачивать по DHT. Это ПРОТИВОРЕЧИТ действующему правилу спеки + ingest «при дедупликации байты сохраняться SHALL NOT» — правило смягчается + этим целевым исключением (MODIFIED-дельта). + +Схема БД не меняется: таблица `download_torrent` уже есть, `source_type` — +существующая колонка; апгрейд — изменение данных, не структуры. + +## Capabilities + +### Modified Capabilities +- `ingest`: уточняется поведение дедупа (пер-хеш гард дозаписи; сохранение + `.torrent`-байт и смена источника при апгрейде `catched`-magnet). + +## Impact + +- Код: `internal/store/download.go` (гард F1 + новый guarded-метод апгрейда), + `internal/ingest/ingest.go` (вызов апгрейда на обоих дедуп-путях). +- Тесты: `internal/store`, `internal/ingest`. +- Совместимость: изменение только ужесточает инвариант (F1) и добавляет + целевой апгрейд (F6); существующие потоки без `.torrent`-поверх-magnet + ведут себя как прежде. diff --git a/openspec/changes/ingest-dedup-integrity/specs/ingest/spec.md b/openspec/changes/ingest-dedup-integrity/specs/ingest/spec.md new file mode 100644 index 0000000..fa66dd5 --- /dev/null +++ b/openspec/changes/ingest-dedup-integrity/specs/ingest/spec.md @@ -0,0 +1,154 @@ +## MODIFIED Requirements + +### Requirement: Атомарность возврата загрузки в активное состояние + +Система SHALL атомарно (в одной write-транзакции) проверять на каждом пути, +возвращающем загрузку из терминального состояния в активное (ручной retry, +воскрешение фоновой сверкой, повторная раскладка/relink) или создающем её +(приём, adopt чужой раздачи), что никакая другая активная загрузка не +владеет любым из хешей этой, и при владении SHALL отказывать в переходе, +сохраняя инвариант «не более одной активной загрузки на infohash». +Отказ SHALL происходить до побочных эффектов во внешних системах +(повторного добавления торрента в qBittorrent). + +Та же проверка SHALL применяться к дозаписи хешей загрузке (раскрытие +гибридного торрента) на ВСЕХ путях дозаписи, включая дедуп-дозапись при +приёме: хеш, которым владеет другая активная загрузка, дописан быть SHALL +NOT — ни отдельным методом дозаписи, ни дедуп-веткой атомарного заведения, +которая доносит недостающие хеши найденной активной задаче. Прямой перевод +терминальной загрузки в активное состояние в обход этой проверки SHALL +отклоняться хранилищем (механический бэкстоп вместо удалённого +unique-индекса). + +#### Scenario: Retry при занятом хеше + +- **GIVEN** загрузка #1 в `failed` с хешем `h`, и другая активная загрузка + #2 с тем же `h` +- **WHEN** пользователь вызывает retry для #1 +- **THEN** переход отклоняется с пояснением, #1 остаётся в `failed` +- **AND** активной по `h` остаётся #2 + +#### Scenario: Дедуп-дозапись не крадёт чужой хеш + +- **GIVEN** активная загрузка A владеет хешем `v1`, активная загрузка B + владеет хешем `v2` того же гибридного торрента +- **WHEN** принимается источник с обоими хешами `{v1, v2}` и дедупится на B +- **THEN** B получает только незанятые хеши, а `v1` (в собственности A) B не + дописывается +- **AND** инвариант «не более одной активной загрузки на infohash» + сохраняется (по `v1` активна только A) + +### Requirement: Приём источника из .torrent-файла + +Приём SHALL принимать источник в виде **байтов `.torrent`-файла** (наряду с +magnet-ссылкой) — тем же быстрым use-case, общим для транспортов. Получив +непустые байты торрента, система SHALL разобрать их локально (без сети), +извлечь инфохэш(и) и завести загрузку с `source_type = torrent`, после чего +сразу вернуть ответ транспорту (синхронный путь к qBittorrent не обращается — +добавление делает воркер, см. `download-tracking`). + +Инфохэши система SHALL извлекать такими, какими их сообщает qBittorrent, чтобы +сопоставление раздач и дедупликация работали: v1-хеш (для v1/гибридного файла) +SHALL вычисляться как SHA1 **исходных** байтов info-словаря (без переэнкода); +v2-хеш (для v2/гибридного файла, BEP52) SHALL извлекаться как 64-hex `infohash_v2`. +Для чистого v2-only файла система SHALL записывать v2-хеш (v1 у него нет). +Извлечение всех известных хешей и дозапись недостающих подчиняются требованию +«Множество инфохэшей загрузки». + +Дедупликацию по активной задаче, атомарное заведение (`download` в состоянии +`catched` + записи `download_infohash`) и инвариант «не более одной активной +загрузки на infohash» torrent-приём SHALL проходить тем же атомарным путём, что +и magnet (см. «Приём источника и заведение загрузки», «Дедупликация приёма по +любому из хешей», «Атомарность возврата загрузки в активное состояние»). + +Байты `.torrent` система SHALL сохранять персистентно, привязанными к загрузке, +чтобы воркер мог добавить источник в qBittorrent именно файлом (не по magnet): +раздачи закрытых трекеров и торренты без DHT по magnet-хешу метаданные не +получат. Сохранение байтов SHALL выполняться в той же write-транзакции, что и +заведение загрузки; при дедупликации (новая загрузка не создана) байты в общем +случае сохраняться SHALL NOT. + +**Исключение — апгрейд пойманной magnet-задачи до torrent.** Если входящий +источник — байты `.torrent`, а дедуп попал на активную загрузку с +`source_type = magnet`, ещё НЕ отданную в qBittorrent (состояние `catched`), +система SHALL в одной write-транзакции сохранить байты `.torrent`, +привязав их к этой загрузке, и сменить её `source_type` на `torrent`. Тем +самым воркер добавит раздачу файлом, а не magnet-хешем (иначе на закрытом +трекере без DHT метаданные не докачаются, а magnet застрянет в metaDL → +failed). Апгрейд SHALL применяться ТОЛЬКО пока загрузка в `catched` (воркер +источник ещё не добавил); для уже добавленной (`downloading` и далее) +загрузки смена `source_type` при дедупе выполняться SHALL NOT — её судьба +решается путями retry/сверки, а не приёмом. Апгрейд SHALL быть best-effort: +его неуспех приём не прерывает. + +Из полей `.torrent` система SHALL синтезировать контекст распознавания (имя +раздачи, суммарный размер, сигнал по дереву файлов, домен трекера, комментарий) +и **дополнять** им контекст транспорта — тем же правилом слияния, что и синтез +из полей magnet (пользовательский текст первым; при пустом тексте — только +синтез). Обогащённый контекст система SHALL сохранять в `download.Context`. +Синтез SHALL выполняться без сетевых запросов. + +`source_ref` у torrent-загрузки SHALL быть человекочитаемым референсом (имя +раздачи или файла), а НЕ адресом добавления: добавление в qBittorrent идёт +байтами, и трактовать `source_ref` как magnet/URL для добавления система SHALL +NOT. + +#### Scenario: Быстрый приём .torrent-файла + +- **GIVEN** валидные байты `.torrent`-файла и (опц.) текст контекста +- **WHEN** вызывается приём +- **THEN** из файла извлекаются инфохэши и создаётся `download` в состоянии + `catched` (`source_type = torrent`) с записями `download_infohash` +- **AND** байты файла сохраняются привязанными к загрузке +- **AND** ответ транспорту отдан без обращения к qBittorrent + +#### Scenario: Инфохэш из исходных байтов info + +- **WHEN** система разбирает v1/гибридный `.torrent`-файл +- **THEN** инфохэш v1 вычисляется как SHA1 исходных байтов info-словаря +- **AND** совпадает с хешем, по которому qBittorrent позже сопоставит раздачу + +#### Scenario: v2-only файл записывается под v2-хешем + +- **WHEN** система разбирает `.torrent` только с метаданными v2 (без v1) +- **THEN** у загрузки записывается v2-хеш (64-hex), совпадающий с `infohash_v2` + qBittorrent +- **AND** сопоставление раздачи работает по нему + +#### Scenario: Дубль .torrent по активной torrent-задаче + +- **GIVEN** уже есть активная (в т.ч. `catched`) загрузка с тем же infohash и + `source_type = torrent` +- **WHEN** принимается `.torrent` с тем же инфохэшем +- **THEN** новая загрузка не создаётся, возвращается существующая +- **AND** байты торрента повторно не сохраняются (дубль) + +#### Scenario: Апгрейд catched-magnet до torrent + +- **GIVEN** активная загрузка в `catched` с `source_type = magnet` и хешем `h` + (magnet-задача ещё не отдана в qBittorrent) +- **WHEN** принимается `.torrent` с тем же инфохэшем `h` +- **THEN** новая загрузка не создаётся, возвращается существующая +- **AND** байты `.torrent` сохраняются привязанными к ней, а её `source_type` + становится `torrent` — в одной транзакции +- **AND** воркер добавит раздачу файлом (не по magnet) + +#### Scenario: Magnet-задача уже добавлена — апгрейда нет + +- **GIVEN** активная загрузка с `source_type = magnet` уже в `downloading` + (отдана в qBittorrent) +- **WHEN** принимается `.torrent` с тем же инфохэшем +- **THEN** возвращается существующая загрузка, её `source_type` остаётся + `magnet`, байты `.torrent` не сохраняются + +#### Scenario: Контекст из полей файла + +- **WHEN** принят `.torrent` с именем раздачи, деревом файлов и трекерами +- **THEN** в `download.Context` добавляется синтез (имя, размер, сигнал по + файлам, домен трекера), дополняющий текст транспорта +- **AND** синтез выполнен без сетевых запросов + +#### Scenario: Слишком большой .torrent отклоняется + +- **WHEN** принимаемый `.torrent`-файл превышает ограничение размера +- **THEN** приём отклоняется с ошибкой, загрузка не создаётся diff --git a/openspec/changes/ingest-dedup-integrity/tasks.md b/openspec/changes/ingest-dedup-integrity/tasks.md new file mode 100644 index 0000000..1ca7ff0 --- /dev/null +++ b/openspec/changes/ingest-dedup-integrity/tasks.md @@ -0,0 +1,29 @@ +## 1. F1 — пер-хеш гард дедуп-дозаписи + +- [x] 1.1 В `CreateDownloadIfNoActive` (`internal/store/download.go`) заменить + безусловный цикл `INSERT OR IGNORE` дедуп-ветки на пер-хеш гард как в + `AddInfohashes`: пропускать хеш, которым владеет другая активная задача + (`findActiveByInfohash(..., existing.ID)`), внутри уже открытой tx. +- [x] 1.2 Тест в `internal/store`: гибридный дедуп на B не крадёт хеш, + принадлежащий активной A (инвариант сохранён). + +## 2. F6 — апгрейд catched-magnet до torrent + +- [x] 2.1 Добавить guarded-метод хранилища + `UpgradeCatchedMagnetToTorrent(ctx, downloadID, torrentBlob) (bool, error)`: + в одной tx гардированным UPDATE `source_type='torrent'` при + `source_type='magnet' AND state='catched'`, затем сохранить байты в + `download_torrent`; вернуть, был ли апгрейд. +- [x] 2.2 В `internal/ingest/ingest.go` свести оба дедуп-пути к `attached()` и + вызвать апгрейд, когда входящий источник — torrent с байтами + (best-effort: неуспех логируется, приём не валится). +- [x] 2.3 Обновить интерфейс `ingest.Store` и `fakeStore` в тестах. +- [x] 2.4 Тесты в `internal/store`: апгрейд из `catched`+magnet сохраняет байты + и меняет `source_type`; из `downloading`/не-magnet — no-op. Тест в + `internal/ingest`: дедуп `.torrent` на catched-magnet вызывает апгрейд. + +## 3. Проверки + +- [x] 3.1 `openspec validate --strict ingest-dedup-integrity` +- [x] 3.2 `task test` +- [x] 3.3 `task lint`