Files
jellybit/internal/ingest/ingest_test.go
T
avandClaude Fable 5 37f2f6481a Идентичность на ULID: download_infohash, guarded-дедуп, миграция (ulid-identity)
Все сущности переехали с INTEGER AUTOINCREMENT на TEXT ULID (lowercase,
internal/ident — единая точка генерации и разбора; oklog/ulid). Инфохэши
загрузки — множество (download_infohash, v1/v2 гибридных торрентов): дедуп
и сопоставление в поллинге по любому из хешей, magnet-парсер отдаёт оба
хеша гибридной ссылки, усечённый v2-хеш v2-only раздач не хранится.

Инвариант «не более одной активной загрузки на infohash» вместо снятого
unique-индекса держат guarded-методы store в одной write-транзакции
(_txlock=immediate): CreateDownloadIfNoActive (приём/adopt, с доносом
недостающих хешей), ActivateIfNoOtherActive (retry/recovery/relink, отказ
до побочных эффектов), guarded AddInfohashes; SetDownloadState отклоняет
терминал→активное как механический бэкстоп.

Миграция 0006 — первая Go-миграция goose: пересоздание таблиц при
включённых FK, backfill ULID с timestamp из created_at (хронология id
сохранена), разнос infohash, удаление idempotency_key. BREAKING: формат id
в URL/логах/Telegram, REST-поля id (string) и infohashes (список).

Новая конвенция docs/conventions/database.md (без числовых PK), корреляция
в логах grep'ом по голому ULID, ER-схема обновлена. Спеки: новая capability
identity, MODIFIED в state-reconciliation; change заархивирован. Пройдены
ревью дизайна и кода (по 8 углов), все находки исправлены с
регрессионными тестами.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-02 21:25:00 +03:00

252 lines
8.1 KiB
Go

package ingest
import (
"context"
"errors"
"io"
"log/slog"
"testing"
"time"
"git.vakhrushev.me/av/jellybit/internal/ident"
"git.vakhrushev.me/av/jellybit/internal/qbt"
"git.vakhrushev.me/av/jellybit/internal/store"
)
const sampleMagnet = "magnet:?xt=urn:btih:541ADCFF3B6DD5DBA7088EA83317D9D6FAC331D6&dn=Dune"
const sampleInfohash = "541adcff3b6dd5dba7088ea83317d9d6fac331d6"
type fakeStore struct {
active *store.Download
created []store.Download
hashes [][]string
toppedUp []string
stateCalls []stateCall
}
type stateCall struct {
id string
state store.State
code string
msg string
}
func (f *fakeStore) FindActiveByInfohash(_ context.Context, _ ...string) (*store.Download, error) {
return f.active, nil
}
func (f *fakeStore) CreateDownloadIfNoActive(_ context.Context, d *store.Download, hashes []string) (*store.Download, error) {
if f.active != nil {
return f.active, nil
}
d.ID = ident.NewID()
f.created = append(f.created, *d)
f.hashes = append(f.hashes, hashes)
return nil, nil
}
func (f *fakeStore) AddInfohashes(_ context.Context, id string, hashes []string) error {
f.toppedUp = append(f.toppedUp, hashes...)
_ = id
return nil
}
func (f *fakeStore) SetDownloadState(_ context.Context, id string, st store.State, code, msg string) error {
f.stateCalls = append(f.stateCalls, stateCall{id, st, code, msg})
return nil
}
type fakeQbt struct {
added []qbt.AddRequest
err error
}
func (f *fakeQbt) Add(_ context.Context, ar qbt.AddRequest) error {
if f.err != nil {
return f.err
}
f.added = append(f.added, ar)
return nil
}
// fakeNamer возвращает заранее заданное имя; фиксирует переданные аргументы.
type fakeNamer struct {
name string
gotContext string
gotHint string
called bool
}
func (f *fakeNamer) DeriveName(_ context.Context, contextText, hint string) string {
f.called = true
f.gotContext = contextText
f.gotHint = hint
return f.name
}
func newService(st Store, qb QBittorrent) *Service {
return newServiceWithNamer(st, qb, nil)
}
func newServiceWithNamer(st Store, qb QBittorrent, nm Namer) *Service {
return New(st, qb, nm, Config{Category: "jellybit", SavePath: "/srv/media/downloads"},
slog.New(slog.NewTextHandler(io.Discard, nil)))
}
func TestIngestHappyPath(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{}
res, err := newService(fs, fq).Ingest(context.Background(), Request{Source: sampleMagnet, Context: "Дюна 2"})
if err != nil {
t.Fatalf("Ingest: %v", err)
}
if len(res.Infohashes) != 1 || res.Infohashes[0] != sampleInfohash {
t.Errorf("infohashes = %v", res.Infohashes)
}
if res.State != store.StateDownloading || res.Deduplicated {
t.Errorf("res = %+v", res)
}
if len(fs.created) != 1 {
t.Fatalf("создано задач: %d, want 1", len(fs.created))
}
if got := fs.created[0]; got.Context != "Дюна 2" {
t.Errorf("сохранённая задача: %+v", got)
}
if len(fs.hashes) != 1 || len(fs.hashes[0]) != 1 || fs.hashes[0][0] != sampleInfohash {
t.Errorf("хеши задачи: %v", fs.hashes)
}
if len(fq.added) != 1 {
t.Fatalf("вызовов qbt.Add: %d, want 1", len(fq.added))
}
add := fq.added[0]
if len(add.URLs) != 1 || add.URLs[0] != sampleMagnet {
t.Errorf("URLs = %v", add.URLs)
}
if add.Category != "jellybit" || add.SavePath != "/srv/media/downloads" {
t.Errorf("category/savepath = %q/%q", add.Category, add.SavePath)
}
}
func TestIngestSetsDisplayName(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{}
nm := &fakeNamer{name: "Дюна: Часть вторая (2024)"}
_, err := newServiceWithNamer(fs, fq, nm).Ingest(context.Background(),
Request{Source: sampleMagnet, Context: "Дюна 2"})
if err != nil {
t.Fatalf("Ingest: %v", err)
}
if !nm.called || nm.gotContext != "Дюна 2" || nm.gotHint != "Dune" {
t.Errorf("namer получил context=%q hint=%q (called=%v)", nm.gotContext, nm.gotHint, nm.called)
}
if len(fq.added) != 1 || fq.added[0].Rename != "Дюна: Часть вторая (2024)" {
t.Errorf("rename = %q, want %q", fq.added[0].Rename, "Дюна: Часть вторая (2024)")
}
// То же имя сохраняется у загрузки — заголовок в веб-UI.
if len(fs.created) != 1 || fs.created[0].DisplayName != "Дюна: Часть вторая (2024)" {
t.Errorf("display_name = %q, want %q", fs.created[0].DisplayName, "Дюна: Часть вторая (2024)")
}
}
func TestIngestEmptyNameOmitsRename(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{}
nm := &fakeNamer{name: ""} // имя не получено
if _, err := newServiceWithNamer(fs, fq, nm).Ingest(context.Background(),
Request{Source: sampleMagnet}); err != nil {
t.Fatalf("Ingest: %v", err)
}
if len(fq.added) != 1 || fq.added[0].Rename != "" {
t.Errorf("rename = %q, want пусто", fq.added[0].Rename)
}
}
func TestIngestIdempotent(t *testing.T) {
existing := &store.Download{ID: "01hzzzexisting000000000000", State: store.StateDownloading}
fs := &fakeStore{active: existing}
fq := &fakeQbt{}
res, err := newService(fs, fq).Ingest(context.Background(), Request{Source: sampleMagnet})
if err != nil {
t.Fatalf("Ingest: %v", err)
}
if !res.Deduplicated || res.DownloadID != existing.ID {
t.Errorf("ожидалось присоединение к существующей задаче: %+v", res)
}
if len(fs.created) != 0 {
t.Error("не должно создаваться новой задачи")
}
if len(fq.added) != 0 {
t.Error("не должно быть повторного добавления в qBittorrent")
}
}
// Быстрый дедуп-путь доносит существующей задаче недостающие хеши
// гибридного magnet (иначе последующий приём по второму хешу создал бы
// вторую активную задачу).
func TestIngestDedupTopsUpHashes(t *testing.T) {
const v2 = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"
existing := &store.Download{
ID: "01hzzzexisting000000000000", State: store.StateDownloading,
Infohashes: []store.Infohash{{DownloadID: "01hzzzexisting000000000000", Infohash: sampleInfohash, Kind: store.HashV1}},
}
fs := &fakeStore{active: existing}
res, err := newService(fs, &fakeQbt{}).Ingest(context.Background(),
Request{Source: sampleMagnet + "&xt=urn:btmh:1220" + v2})
if err != nil {
t.Fatalf("Ingest: %v", err)
}
if !res.Deduplicated {
t.Fatalf("ожидался дедуп: %+v", res)
}
if len(fs.toppedUp) != 2 {
t.Errorf("хеши не донесены существующей задаче: %v", fs.toppedUp)
}
}
func TestIngestQbitErrorMarksFailed(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{err: errors.New("connection refused")}
res, err := newService(fs, fq).Ingest(context.Background(), Request{Source: sampleMagnet})
if err == nil {
t.Fatal("ожидалась ошибка")
}
if res.State != store.StateFailed {
t.Errorf("state = %q, want failed", res.State)
}
if len(fs.stateCalls) != 1 || fs.stateCalls[0].state != store.StateFailed {
t.Errorf("ожидался перевод в failed: %+v", fs.stateCalls)
}
}
func TestIngestQbitErrorNotifies(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{err: errors.New("connection refused")}
svc := newService(fs, fq)
got := make(chan string, 1)
svc.SetFailureNotifier(func(id string) { got <- id })
if _, err := svc.Ingest(context.Background(), Request{Source: sampleMagnet}); err == nil {
t.Fatal("ожидалась ошибка")
}
select {
case id := <-got:
if id == "" {
t.Errorf("уведомление с пустым id")
}
case <-time.After(2 * time.Second):
t.Fatal("уведомление о падении приёма не пришло")
}
}
func TestIngestRejectsNonMagnet(t *testing.T) {
fs := &fakeStore{}
fq := &fakeQbt{}
if _, err := newService(fs, fq).Ingest(context.Background(), Request{Source: "https://example.com/x.torrent"}); err == nil {
t.Fatal("ожидалась ошибка для не-magnet источника")
}
if len(fs.created) != 0 || len(fq.added) != 0 {
t.Error("не должно быть ни записи, ни добавления")
}
}